Skip to main content
Reliability, async & real-time
決済
信頼性
アーキテクチャ設計
PostgreSQL
Python

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
Reading time
20 min read
Author
友田 陽大
Share

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). 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.


Double charges, missed webhooks, or churn from failed payments?

Payment-reliability diagnosis and billing architecture, as a technical advisor

2. Separate three kinds of state

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

StateMeaningWhere it lives in this article's example
DesiredWhat your business has decided it should bepayments.status (pending / paid / failed)
Known appliedThe external state you last confirmed, and whenpayments.observed_external_status / last_reconciled_at
External authoritativeWhat the external system holds now, treated as the truthThe 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:

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:

# ❌ 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:

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

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

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.

5.3 Running one batch

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.

-- 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.

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). Delivering an outbox with several workers is covered in leases and fencing tokens in an outbox dispatcher.

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).

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.

TestExpected result
Classification tablepending×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 row10 UNVERIFIABLE, 0 IN_SYNC, no last-verified times advanced
Detect drift over two cyclesOne finding and one command per kind of drift
paid locally, canceled externallyA NEEDS_REVIEW finding only; no repair command
Row updated within the grace periodIN_FLIGHT; the external API isn't called
Another reconciler advanced the cursor firstConcurrentRun; the batch is discarded
Repair races a webhookOnly whichever applies first changes state; the other is a no-op
paid locally, requires_action externallyNEEDS_REVIEW, not IN_FLIGHT
Applying a repair commandThe finding is resolved, so the same drift can be detected again
Rate limiter11 calls at 50 requests per second take at least 0.19 s
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. Stripe webhooks and the subscription state machine are covered in Implementing Stripe webhooks and idempotency to production quality.

Frequently asked questions

What is reconciliation?
A mechanism that periodically compares the state in your database with the state held by the external system of record, such as Stripe, detects drift and repairs it. It is a safety net that converges drift caused by missed webhooks or outages after the fact; it does not replace transactions or idempotency.
Can drift found by reconciliation be fixed automatically?
It depends on the direction. Bringing your state up to date with the external system — say, an order still marked unpaid locally because a webhook was missed — is easy to automate if it goes through the canonical write path. Drift where you shipped goods as paid but the external payment never succeeded involves decisions about refunds or stopping shipments, so send it to a human instead.
If I reconcile LIMIT 200 every hour, will I eventually check every row?
No. LIMIT on an ordering such as ORDER BY updated_at DESC picks nearly the same rows each time, and old rows are never checked. Persist a keyset cursor that advances on a key such as the primary key, record when each cycle completes and when each row was last verified, and measure that every row has been checked within a bounded time.
How should a reconciler treat rows it couldn't check because the external API failed?
Don't count them as healthy; count them as unverifiable, and don't advance their last-verified time. That way a long external outage doesn't produce a false 'zero drift', and unchecked rows show up in monitoring through the age of their last verification.

References

友田

友田 陽大

Developer of a METI Minister's Award–winning product. With TypeScript + Python + AWS, I deliver SaaS, industry DX, and production-grade generative AI (RAG) end to end — from requirements to infrastructure and operations — single-handedly.

Double charges, missed webhooks, or churn from failed payments?

Payment-reliability diagnosis and billing architecture, as a technical advisor

"We get a double charge now and then." "We find out about dropped webhooks after the fact." "Failed payments keep turning into cancellations." The cause is rarely one bug — it is usually a structural gap in idempotency, ordering, retry policy, or the state machine. From reading your existing implementation to a receiver that survives redelivery, a recovery flow for failed payments, and where to draw the build-vs-buy line, we decide it together.

Available for both project-based (contract) and advisory engagements. Start with a free 30-minute consult.

Also worth reading