# Designing Reconciliation: Detect and Classify Drift Between Your Database and Stripe, Then Repair It Safely Through an Outbox

> When missed webhooks or outages leave your database and an external service such as Stripe out of step, reconciliation detects the drift, classifies it, and sends repair commands through the canonical write path. Why 'LIMIT 200 every hour' never covers everything, a keyset cursor with a full-cycle guarantee, classification that never treats unverifiable state as healthy, rate limits and monitoring — with Python code tested against PostgreSQL 18.

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

## Key points

- Even with correct webhooks, outboxes and idempotency, your database and an external service drift apart because of events never sent, long outages, bugs and manual operations. Reconciliation finds and fixes that drift after the fact; it does not replace transactions
- Separate three states: what your business intends (desired), what you last confirmed externally (known applied), and what the external system holds now (external authoritative). Write drift detection and repair as comparisons between them
- 'Run it hourly with LIMIT 200 and it'll eventually see everything' is false. ORDER BY updated_at DESC LIMIT 200 sees the same rows every time. Advance a keyset cursor, and measure coverage with completed cycles and a per-row last-verified time
- An external API failure (timeout, 429) means 'unverifiable', not 'healthy'. Advance the last-verified time only for rows you actually confirmed, and monitor the unverifiable rate
- The reconciler doesn't UPDATE directly. It queues a repair command in an outbox, applied by the same canonical write path as webhooks. Drift involving money or shipped goods is not auto-repaired; it goes to a human

---

**The short answer.** Even with idempotent webhook handling, a transactional outbox and idempotency keys all implemented correctly, your database and an external service such as Stripe will drift apart eventually — because of events the provider never sent, outages longer than the webhook retry window, bugs, and manual changes in a dashboard. **Reconciliation** periodically detects that drift, classifies it and repairs it. Four things matter in the design: **compare three distinct kinds of state**, **measure the guarantee that every row is checked within a bounded time**, **never treat state you couldn't verify as healthy**, and **don't let the reconciler write repairs itself — hand them as commands to the same canonical write path that webhooks use.**

Reconciliation is not a **replacement** for transactions or idempotency. Prevent with atomicity whatever atomicity can prevent; reconciliation is the **last safety net** that converges whatever drift remains.

The code was **verified on PostgreSQL 18.6 / psycopg 3.3 with 12 tests, including a reproduction of "LIMIT 200 never sees old rows"** (October 2026).

> **Assumption in the code**: the psycopg code assumes connections opened with `psycopg.connect(dsn, autocommit=True)` and expresses transaction boundaries only with `with conn.transaction():`. Put `from __future__ import annotations` at the top of the module: `run_batch` is annotated with `RateLimiter`, which is defined later in section 9.

---

## 1. What kinds of drift happen

Two concrete examples.

**Example 1: unpaid locally, paid in Stripe.** The endpoint was down for a long time before the `payment_intent.succeeded` webhook arrived, and Stripe's retry window (up to three days in live mode) passed. Or the webhook handler had a "record processed first" bug and the event was lost (Case A in [implementing webhook idempotency with a transactional inbox](/blog/webhook-idempotency-transactional-inbox-postgresql-guide)). The customer paid, but the order is stuck at "awaiting payment".

**Example 2: paid locally, payment cancelled in Stripe.** You marked the order `paid` and shipped it, but the PaymentIntent was later cancelled in Stripe — because of a missed event, a manual change in the dashboard, an API call from another system, and so on.

Example 1 is drift you fix by "bringing yourself up to date with the external system". Example 2 involves **decisions about money and shipped goods**; it isn't solved by mechanically reverting `paid`. Reconciliation design starts by telling these two apart.

---

## 2. Separate three kinds of state

To handle drift precisely, don't represent "the state" as one column; split it into three.

| State | Meaning | Where it lives in this article's example |
| --- | --- | --- |
| **Desired** | What your business has decided it should be | `payments.status` (`pending` / `paid` / `failed`) |
| **Known applied** | The external state you last confirmed, and when | `payments.observed_external_status` / `last_reconciled_at` |
| **External authoritative** | What the external system holds now, treated as the truth | The PaymentIntent `status` from the Stripe API |

There are two reasons to keep the known-applied state. The first is **monitoring**: a row with an old `last_reconciled_at` is one nobody has compared with the outside for a while, so you can measure whether the reconciler really covers everything (section 5). The second is **change detection**: whether the external state changed since your last check can set an investigation's priority.

Which side is "the truth" is decided per kind of data. For whether a payment succeeded, Stripe is the truth. For "should this customer get a discount", your business decision is the truth and Stripe is where it's applied. In the first case you fix your side to match the external system; in the second you fix the external system to match yours. **Write "if they differ, fix it" without deciding which side is the truth, and the two systems keep overwriting each other.**

---

## 3. What the reconciler does, and doesn't do

Reconciliation implementations tend to fail at one of two extremes:

- **SELECT and log**: the drift is found, buried in logs nobody reads, and never fixed.
- **A second business implementation that repairs everything directly**: the reconciler writes `UPDATE payments SET status = 'paid'` itself. Now the same business logic as the webhook handler (booking the ledger, notifications, confirming stock) has to live in the reconciler too, and the two implementations slowly diverge. Run alongside a webhook, it also double-books.

The flow I recommend:

```text
detect (compare with the external system)
   ↓
classify (decide the kind of drift and whether it may be auto-repaired; a pure function)
   ↓
record finding (record the drift; never register the same drift twice)
   ↓
repair command (queue a repair command in the outbox, only for auto-repairable drift)
   ↓
canonical writer (applied by the same canonical write path as webhooks; an idempotent state transition)
```

The reconciler's responsibility ends at "find drift and request a repair". **The actual state change is made only by the single canonical write path that the webhook handler also uses.** That keeps business logic in one place, and the idempotent state transition stays safe when a webhook and a repair run at the same time.

---

## 4. ❌ Why "LIMIT 200 every hour" never covers everything

This is the implementation you see most:

```python
# ❌ Broken: "LIMIT 200 every hour" does not guarantee every row is seen. Do not copy this
def scan_broken(conn: psycopg.Connection, gateway: PaymentGateway, *, batch: int = 200) -> int:
    rows = conn.execute(
        """
        SELECT id, stripe_payment_intent_id FROM payments
         ORDER BY updated_at DESC
         LIMIT %s
        """,
        (batch,),
    ).fetchall()
    for payment_id, pi_id in rows:
        gateway.payment_intent_status(pi_id)
        conn.execute("UPDATE payments SET last_reconciled_at = now() WHERE id = %s", (payment_id,))
    return len(rows)
```

It's tempting to think that running this hourly will eventually see everything, but **it sees almost the same 200 rows every time.** `ORDER BY updated_at DESC LIMIT 200` picks "the 200 most recently updated rows"; old rows sink further as new rows arrive and are never chosen again.

In my tests, running this function ten times (ten hours' worth) over 500 payments left **300 rows that were never checked.** Without `ORDER BY` it's no better: PostgreSQL doesn't guarantee the order of a `LIMIT` without `ORDER BY`, so which rows you get depends on the plan, and "it happens to be the same rows every time" is common.

And it's precisely **old rows** that are most likely to have drifted. Older transactions past the webhook retry window have no rescue other than reconciliation.

### 4.1 How long a full cycle takes

Think about coverage in numbers:

```text
N = rows in scope, B = batch size, R = runs per hour, G = new rows per hour
time for one full cycle ≈ N ÷ (B × R)     and unless B × R > G, a cycle never completes
```

For example, with N = 120,000 rows, B = 200 and R = 12 (every five minutes), a cycle takes 120,000 ÷ 2,400 = 50 hours. You can claim "drift is found within about two days" only when this calculation holds **and you are measuring** that cycles actually complete.

---

## 5. Advance with a keyset cursor and record completed cycles

### 5.1 Schema

```sql
CREATE TABLE payments (
    id                       bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    order_id                 text        NOT NULL,
    stripe_payment_intent_id text        NOT NULL UNIQUE,
    status                   text        NOT NULL CHECK (status IN ('pending', 'paid', 'failed')),
    version                  bigint      NOT NULL DEFAULT 0,
    created_at               timestamptz NOT NULL DEFAULT now(),
    updated_at               timestamptz NOT NULL DEFAULT now(),
    observed_external_status text,         -- known applied: the status last confirmed in Stripe
    last_reconciled_at       timestamptz   -- when this row was last compared externally (advances only on success)
);

CREATE TABLE reconciliation_cursor (
    job                     text        PRIMARY KEY,
    last_id                 bigint      NOT NULL DEFAULT 0,
    cycle_started_at        timestamptz NOT NULL DEFAULT now(),
    last_cycle_completed_at timestamptz
);

CREATE TABLE reconciliation_findings (
    id               bigint GENERATED ALWAYS AS IDENTITY PRIMARY KEY,
    payment_id       bigint      NOT NULL REFERENCES payments (id),
    kind             text        NOT NULL,
    local_status     text        NOT NULL,
    external_status  text,
    detected_at      timestamptz NOT NULL DEFAULT now(),
    resolved_at      timestamptz
);
-- Don't register the same drift on every scan (one open finding at a time)
CREATE UNIQUE INDEX reconciliation_findings_open
    ON reconciliation_findings (payment_id, kind) WHERE resolved_at IS NULL;
```

The cursor advances on the primary key `id`. `WHERE id > last_id ORDER BY id LIMIT B` reads a range of the primary-key index, so the cost per run stays constant as the table grows. `OFFSET` slows down in proportion to the rows it skips, and rows added or removed mid-scan cause skips or double reads.

If `id` isn't monotonic (a UUID primary key, say), use **a unique, ordered pair** such as `(created_at, id)` as the cursor. PostgreSQL's row constructor comparison `WHERE (created_at, id) > (%s, %s) ORDER BY created_at, id` handles this.

### 5.2 Classification: a pure function with no I/O

```python
from dataclasses import dataclass
from enum import Enum


class Verdict(Enum):
    IN_SYNC = "in_sync"
    IN_FLIGHT = "in_flight"            # changed recently / external still processing; don't judge now
    REPAIRABLE = "repairable"          # may send a repair command to the canonical write path
    NEEDS_REVIEW = "needs_review"      # never auto-repair (drift involving money or shipped goods)
    UNVERIFIABLE = "unverifiable"      # couldn't check externally; not considered healthy


# External (Stripe PaymentIntent) statuses that are still in transition
_EXTERNAL_IN_FLIGHT = {"processing", "requires_action", "requires_confirmation", "requires_capture"}


@dataclass(frozen=True)
class Decision:
    verdict: Verdict
    kind: str | None = None            # kind of finding
    repair_command: str | None = None  # command name to queue in the outbox


def classify(local_status: str, external_status: str) -> Decision:
    if local_status == "pending" and external_status in _EXTERNAL_IN_FLIGHT:
        # Neither side has settled yet. paid locally but in transition externally falls through to ("paid", _)
        return Decision(Verdict.IN_FLIGHT)
    match (local_status, external_status):
        case ("paid", "succeeded") | ("failed", "canceled") | ("pending", "requires_payment_method"):
            return Decision(Verdict.IN_SYNC)
        case ("pending", "succeeded"):
            # A missed webhook. Advance to paid through the same canonical write path as the webhook
            return Decision(Verdict.REPAIRABLE, "missed_success", "ApplyPaymentSucceeded")
        case ("pending", "canceled"):
            return Decision(Verdict.REPAIRABLE, "missed_cancel", "ApplyPaymentCanceled")
        case ("paid", _):
            # We treated it as paid (and may have shipped or granted access), but externally it didn't succeed
            return Decision(Verdict.NEEDS_REVIEW, "paid_locally_not_externally")
        case _:
            return Decision(Verdict.NEEDS_REVIEW, f"unexpected:{local_status}->{external_status}")
```

Classification policy:

- **Bringing your side up to date with the external system** (`pending` locally → `succeeded` / `canceled` externally) is `REPAIRABLE`. It just triggers the state transition the webhook would have caused, through the same path.
- **Your side being ahead** (`paid` locally, not successful externally) is `NEEDS_REVIEW`. Stopping shipment, contacting the customer and deciding on a refund are business decisions the reconciler has no business making.
- **`pending` locally with the external system mid-transition** (`processing` and so on) is `IN_FLIGHT`. Judging now would report as drift a state that becomes correct seconds later. **`paid` locally while the external system is mid-transition** (say, `requires_action` abandoned at authentication) is `NEEDS_REVIEW`, not `IN_FLIGHT`. Because the external state *was* read, `IN_FLIGHT` would keep advancing the last-verified time, and "lag monitoring green, payment never actually received" would persist forever as false healthy (verified).
- **Unexpected combinations** fall to `NEEDS_REVIEW`, so nothing unknown is classified as "healthy".

As a pure function, the classifier can be covered exhaustively with table-driven tests and no I/O. Check the full list of PaymentIntent `status` values in the [Stripe API reference](https://docs.stripe.com/api/payment_intents/object).

### 5.3 Running one batch

```python
from dataclasses import field
from datetime import timedelta
from typing import Protocol

import psycopg
from psycopg.types.json import Jsonb


class GatewayUnavailable(Exception):
    """Timeout, 429 or 5xx: the external state could not be checked."""


class ExternalNotFound(Exception):
    pass


class PaymentGateway(Protocol):
    """Boundary around the Stripe SDK. Translating SDK exceptions into the two above is this layer's job."""
    def payment_intent_status(self, payment_intent_id: str) -> str: ...


@dataclass
class BatchReport:
    scanned: int = 0
    counts: dict[Verdict, int] = field(default_factory=lambda: {v: 0 for v in Verdict})
    wrapped: bool = False


class ConcurrentRun(Exception):
    """Another reconciler advanced the cursor first. Discard this batch's writes."""


def run_batch(
    conn: psycopg.Connection,
    gateway: PaymentGateway,
    *,
    job: str = "payments-vs-stripe",
    batch: int = 200,
    grace: timedelta = timedelta(minutes=10),
    limiter: RateLimiter | None = None,
) -> BatchReport:
    report = BatchReport()

    # 1) Read the cursor and target rows in a short transaction (hold no locks during external calls)
    with conn.transaction():
        conn.execute("INSERT INTO reconciliation_cursor (job) VALUES (%s) ON CONFLICT DO NOTHING", (job,))
        cursor_row = conn.execute(
            "SELECT last_id FROM reconciliation_cursor WHERE job = %s", (job,)
        ).fetchone()
        assert cursor_row is not None  # always exists after the INSERT ... ON CONFLICT above
        cursor_id: int = cursor_row[0]
        rows = conn.execute(
            """
            SELECT id, stripe_payment_intent_id, status,
                   updated_at > now() - %s AS recently_changed
              FROM payments
             WHERE id > %s
             ORDER BY id
             LIMIT %s
            """,
            (grace, cursor_id, batch),
        ).fetchall()

    if not rows:
        # Reached the end = one full cycle. Rewind to the start (CAS so concurrent runs can't collide)
        with conn.transaction():
            cur = conn.execute(
                """
                UPDATE reconciliation_cursor
                   SET last_id = 0, last_cycle_completed_at = now(), cycle_started_at = now()
                 WHERE job = %s AND last_id = %s
                """,
                (job, cursor_id),
            )
            if cur.rowcount != 1:
                raise ConcurrentRun(job)
        report.wrapped = True
        return report

    # 2) Read external state outside any transaction
    results: list[tuple[int, str, Decision, str | None]] = []
    for payment_id, pi_id, local_status, recently_changed in rows:
        if recently_changed:
            results.append((payment_id, local_status, Decision(Verdict.IN_FLIGHT), None))
            continue
        if limiter:
            limiter.wait()
        try:
            external = gateway.payment_intent_status(pi_id)
        except GatewayUnavailable:
            results.append((payment_id, local_status, Decision(Verdict.UNVERIFIABLE), None))
            continue
        except ExternalNotFound:
            results.append(
                (payment_id, local_status, Decision(Verdict.NEEDS_REVIEW, "missing_externally"), None)
            )
            continue
        results.append((payment_id, local_status, classify(local_status, external), external))

    # 3) Record results, queue repair commands and advance the cursor in one transaction
    with conn.transaction():
        for payment_id, local_status, decision, observed in results:
            report.counts[decision.verdict] += 1
            if observed is not None:
                # Advance "verified" only for rows actually confirmed externally (not for UNVERIFIABLE)
                conn.execute(
                    """
                    UPDATE payments SET observed_external_status = %s, last_reconciled_at = now()
                     WHERE id = %s
                    """,
                    (observed, payment_id),
                )
            if decision.kind is None:
                continue
            finding = conn.execute(
                """
                INSERT INTO reconciliation_findings (payment_id, kind, local_status, external_status)
                VALUES (%s, %s, %s, %s)
                ON CONFLICT (payment_id, kind) WHERE resolved_at IS NULL DO NOTHING
                RETURNING id
                """,
                (payment_id, decision.kind, local_status, observed),
            ).fetchone()
            if finding is not None and decision.repair_command is not None:
                # No direct UPDATE. Hand a command to the same canonical write path as webhooks
                conn.execute(
                    """
                    INSERT INTO outbox (aggregate_type, aggregate_id, event_type, payload)
                    VALUES ('Payment', %s, %s, %s)
                    """,
                    (
                        str(payment_id),
                        decision.repair_command,
                        Jsonb({"payment_id": payment_id, "finding_id": finding[0],
                               "expected_local_status": local_status}),
                    ),
                )
        cur = conn.execute(
            "UPDATE reconciliation_cursor SET last_id = %s WHERE job = %s AND last_id = %s",
            (rows[-1][0], job, cursor_id),
        )
        if cur.rowcount != 1:
            raise ConcurrentRun(job)  # the exception discards the whole transaction
    report.scanned = len(rows)
    return report
```

The design decisions in this code:

1. **External APIs are called outside any transaction.** No connection or lock is held during 200 × a few hundred milliseconds of API calls.
2. **Recording results and advancing the cursor happen in one transaction.** If the process dies midway, neither the cursor nor the records move, and the next run redoes the same batch. Findings can't duplicate thanks to the partial unique index, so redoing is safe.
3. **The cursor advances with CAS (compare-and-set).** It updates `WHERE last_id = <the value read>`, and zero rows means another reconciler got there first, so the batch is discarded. Even with duplicate runs, records and commands aren't doubled (verified).
4. **Rows updated within the grace period (`grace`) aren't judged.** A row a webhook just updated can look like drift because of propagation delay between the two systems; it's checked on the next cycle.
5. **Repair commands are queued only when a new finding is created.** If the same drift is still open on the next cycle, there's still only one command (verified).

### 5.4 Measuring the full-cycle guarantee

When the cursor reaches the end it records `last_cycle_completed_at` and rewinds. In my tests, 450 rows with a batch of 200 advanced 200 → 200 → 50 → 0 (cycle complete, rewind), and every row's `last_reconciled_at` was filled in.

But a completed cursor cycle doesn't mean every row was **successfully checked**: rows the external API failed on were skipped as `UNVERIFIABLE`. So the real guarantee is measured with **each row's last-verified time.**

```sql
-- Reconciliation lag: the age of the row that has gone longest without an external comparison (the key metric)
SELECT max(now() - coalesce(last_reconciled_at, created_at)) AS reconciliation_lag
  FROM payments;

-- Cycle duration and time since the last completed cycle
SELECT job,
       last_cycle_completed_at,
       now() - last_cycle_completed_at AS since_last_cycle,
       now() - cycle_started_at        AS current_cycle_age
  FROM reconciliation_cursor;
```

Put an SLO on `reconciliation_lag` (say, 72 hours) and alert when it's exceeded. This also catches the state where "the reconciler is running but external API errors mean it isn't actually verifying anything."

---

## 6. Don't count "unverifiable" as "healthy"

When the external API times out, swallowing the error in a `try/except` and moving to the next row makes that row look like **nothing happened — no drift.** If Stripe is unreachable for a day, the reconciler reports "zero drift" all day. That's **false healthy.**

In this article's code, rows that hit `GatewayUnavailable` are counted as `UNVERIFIABLE`, and their `last_reconciled_at` is **not advanced.** In my tests, when the external side returned errors for every row, `UNVERIFIABLE` was 10, `IN_SYNC` was 0, and no row's last-verified time moved.

In monitoring, emit the per-batch `UNVERIFIABLE` ratio as a metric, and never judge health from the `IN_SYNC` count alone.

Rows not found externally (404) are classified `NEEDS_REVIEW` (`missing_externally`), not `UNVERIFIABLE`. "We checked, and it doesn't exist" is a clear fact — test data leaking in, a key for another account, a deletion — that a human should investigate.

---

## 7. Repair: hand it to the canonical write path

The `ApplyPaymentSucceeded` command the reconciler queued in the outbox is picked up by the dispatcher and applied with the same function as the webhook handler.

```python
def apply_payment_succeeded(conn: psycopg.Connection, payment_id: int, *, finding_id: int | None = None) -> bool:
    """The pending → paid transition. Does nothing if already advanced (safe against webhook/repair races).

    When called from a repair command, resolve the finding even if the transition is a no-op
    (the target state has been reached either way).
    """
    with conn.transaction():
        cur = conn.execute(
            """
            UPDATE payments SET status = 'paid', version = version + 1, updated_at = now()
             WHERE id = %s AND status = 'pending'
            """,
            (payment_id,),
        )
        if finding_id is not None:
            conn.execute(
                "UPDATE reconciliation_findings SET resolved_at = now() WHERE id = %s AND resolved_at IS NULL",
                (finding_id,),
            )
        return cur.rowcount == 1
```

`WHERE status = 'pending'` is what makes a repair racing a webhook safe. If a late webhook moved the row to `paid` between the reconciler finding the drift and the command being applied, the repair updates zero rows — a no-op (verified). In a real system this function also books the ledger and records notifications in the outbox. **Because webhooks and repairs go through the same function, that logic is written in exactly one place.**

Routing repair commands through the outbox makes the reconciler's transaction (recording the finding) and the repair request atomic. If the reconciler sent directly to the dispatcher's queue or a broker, you'd recreate the two-write problem: "finding recorded, command never sent" ([transactional outbox](/blog/transactional-outbox-pattern-reliable-event-publishing-guide)). Delivering an outbox with several workers is covered in [leases and fencing tokens in an outbox dispatcher](/blog/outbox-dispatcher-skip-locked-lease-fencing-token-guide).

When called from a repair command, pass `finding_id`, and **resolve the finding even if the transition is a no-op.** If a finding stays open, the partial unique index means the same drift can never be registered again, so a repair command that ended up `dead` would never be retried. In operation, alert on findings left open longer than some threshold (say, 24 hours).

`NEEDS_REVIEW` findings go to a human through an admin screen or a ticket. What the human decides is also applied through the canonical write path, and the finding's `resolved_at` is filled in. Keeping a record of who decided what, and when, matters for payment audits.

---

## 8. The other direction: resources that exist only externally

The scan so far runs "from your rows to the external system". In that direction you can't find **resources that exist only externally (orphans)** — for instance, a PaymentIntent created in Stripe just before your own transaction rolled back.

To find them, run a separate job in the opposite direction: page through the external system's list and check each item against your database. Stripe's list APIs use cursor-based pagination with `starting_after` (an object ID), returning up to 100 items at a time, newest first. Bound the period each job covers (say, objects created in the last 7 days), persist the page cursor, and advance a little at a time in the same way.

You can also search `metadata` with Stripe's Search API, but the documentation says **data is normally searchable in under a minute and can be delayed during an outage**, and warns against using search for read-after-write flows. It's usable for reconciliation, where you compare after some time has passed, but don't conclude that "not found" means "doesn't exist".

---

## 9. Rate limits: don't eat your production API budget

A reconciler calls external APIs a lot. According to Stripe's documentation, rate limits are per Stripe account: the live-mode global limit is **100 requests per second**, and individual endpoints are **25 requests per second** unless otherwise noted (as of October 2026; limits change, so check the documentation).

**That budget is shared with your payment traffic itself.** If the reconciler uses all of it, customers' payments fail with 429. Give the reconciler only a small fraction (for example, 5–10 requests per second).

```python
import time


class RateLimiter:
    """A simple fixed-interval limiter. Uses only a fraction of the production API rate limit."""

    def __init__(self, rps: float) -> None:
        self._interval = 1.0 / rps
        self._next = time.monotonic()

    def wait(self) -> None:
        now = time.monotonic()
        if now < self._next:
            time.sleep(self._next - now)
        self._next = max(now, self._next) + self._interval
```

That allocation feeds directly into the calculation in section 4.1. At 5 requests per second you get 18,000 checks an hour, so a 120,000-row cycle takes about seven hours. When growth makes the SLO unattainable, consider **narrowing what you compare** (for example, recent 90-day transactions and unfinished states frequently, everything else rarely) before raising the rate.

On a 429, treat it as `GatewayUnavailable`; stopping the rest of the batch's calls and leaving them for the next run is the safe choice (this article's code decides per row for simplicity).

---

## 10. Tests

These tests were run against a PostgreSQL 18 container and all pass. The external API is replaced by a fake that returns statuses from a dictionary.

| Test | Expected result |
| --- | --- |
| Classification table | `pending`×`succeeded` → `REPAIRABLE`, `paid`×`canceled` → `NEEDS_REVIEW`, `processing` → `IN_FLIGHT` |
| ❌ `ORDER BY updated_at DESC LIMIT 200` ten times (500 rows) | 300 rows never checked (failure reproduced) |
| Keyset cursor (450 rows, batch 200) | 200 → 200 → 50 → 0 (cycle complete); every row's last-verified time filled |
| External API fails for every row | 10 `UNVERIFIABLE`, 0 `IN_SYNC`, no last-verified times advanced |
| Detect drift over two cycles | One finding and one command per kind of drift |
| `paid` locally, `canceled` externally | A `NEEDS_REVIEW` finding only; no repair command |
| Row updated within the grace period | `IN_FLIGHT`; the external API isn't called |
| Another reconciler advanced the cursor first | `ConcurrentRun`; the batch is discarded |
| Repair races a webhook | Only whichever applies first changes state; the other is a no-op |
| `paid` locally, `requires_action` externally | `NEEDS_REVIEW`, not `IN_FLIGHT` |
| Applying a repair command | The finding is resolved, so the same drift can be detected again |
| Rate limiter | 11 calls at 50 requests per second take at least 0.19 s |

```python
def test_broken_scan_starves_old_rows(conn):
    seed(conn, 500)
    gw = FakeStripe({f"pi_{i}": "succeeded" for i in range(1, 501)})
    for _ in range(10):  # ten hours' worth
        scan_broken(conn, gw)
    never = conn.execute("SELECT count(*) FROM payments WHERE last_reconciled_at IS NULL").fetchone()[0]
    assert never == 300  # 300 rows are never seen, even after ten runs
```

---

## 11. Avoiding misuse

- **Don't use reconciliation instead of transactions.** Skip atomicity because "it'll get fixed later" and the volume of drift outgrows the reconciler, and `NEEDS_REVIEW` becomes more than humans can handle.
- **Don't put business logic in the reconciler.** Repairs are commands to the canonical write path.
- **Decide which side is the truth per kind of data before writing any repair.** Fix in both directions without deciding, and the systems overwrite each other.
- **Start automatic repair with "bring yourself up to date with the external system".** The reverse direction, involving money, shipped goods or access, goes through human review.
- **Measure the guarantee in time, not counts.** Monitor "every row verified within 72 hours", not "200 rows per run".

---

## Summary

- Reconciliation is a safety net that converges drift atomicity and idempotency can't prevent; it doesn't replace transactions.
- Split state into desired / known applied / external authoritative, and decide per kind of data which is the truth.
- "LIMIT 200 every hour" keeps seeing the same rows. Advance a keyset cursor and measure coverage with completed cycles and per-row last-verified times.
- An external API failure is "unverifiable", not "healthy". Don't advance the last-verified time, and monitor the ratio.
- Repairs go through the outbox to the canonical write path. Drift involving money or shipped goods goes to a human.
- The external API budget is shared with production payments; give the reconciler only a fraction.

How to converge an external API call whose outcome is unknown into exactly one effect, with idempotency keys and reconciliation, is covered in [What exactly-once actually guarantees](/blog/exactly-once-at-least-once-idempotency-effectively-once-guide). Stripe webhooks and the subscription state machine are covered in [Implementing Stripe webhooks and idempotency to production quality](/blog/stripe-payments-production-guide-webhooks-idempotency-subscriptions).
