Skip to content

[WIP] Add idempotency and resilience to M-Pesa webhook handling - #5

Draft
Victor-Kipruto-Rop with Copilot wants to merge 1 commit into
mainfrom
copilot/add-db-unique-constraint-trans-id
Draft

[WIP] Add idempotency and resilience to M-Pesa webhook handling#5
Victor-Kipruto-Rop with Copilot wants to merge 1 commit into
mainfrom
copilot/add-db-unique-constraint-trans-id

Conversation

Copilot AI commented Sep 3, 2026

Copy link
Copy Markdown
Contributor

Thanks for asking me to work on this. I will get started on it and keep this PR's description up to date as I form a plan and make progress.

Background

The repository currently can receive duplicate M-Pesa webhook deliveries which can create duplicate transaction rows in the database and lead to duplicate downstream processing. There is also no idempotency check before publishing to Kafka, and the consumer may insert duplicates if duplicate messages are consumed. We must implement end-to-end idempotency and resilience:

Goals

  • Add a DB-level unique constraint on the transaction identifier (trans_id) to prevent persistent duplicates.
  • Short-circuit duplicate webhook deliveries using an atomic in-memory/cache check (Redis SET NX).
  • Ensure the webhook handler only publishes to Kafka after a successful DB insert; if Kafka publish fails, mark the row for retry instead of deleting it.
  • Make the consumer idempotent by switching to an INSERT ... ON CONFLICT DO NOTHING pattern.
  • Add tests and dependency updates where needed.

Required changes (high-level)

  1. Database migration
  • Add a new SQL migration file migrations/20260903_add_unique_index_transid.sql that attempts to create a unique index on transactions.trans_id.
  • The migration must perform a pre-check for duplicates and abort (raise) with a helpful message if duplicates exist. It should recommend running a dedupe script or manual cleanup before applying.
  • Prefer creating the unique index CONCURRENTLY (Postgres) to avoid long locks, with guidance in comments that CONCURRENTLY cannot run inside a transaction.
  • Also add an Alembic revision under alembic/versions/ that wraps the same logic and fails early if duplicates exist.
  1. Webhook handler changes
  • Locate the existing webhook handler(s) that process incoming MPESA webhooks (search for keywords: webhook, mpesa, transaction, trans_id, TransactionId, kafka).
  • Replace the current flow (publish to Kafka then DB write or DB write without idempotency) with this safe flow:
    a) Extract and normalize trans_id from the payload.
    b) Use an atomic cache check with Redis: redis.set(key, 1, ex=TTL, nx=True). If not set => return 200 duplicate.
    c) Attempt an atomic DB insert using INSERT ... ON CONFLICT DO NOTHING (or SQLAlchemy's pg_insert(...).on_conflict_do_nothing()). If the insert rowcount==0 => duplicate => return 200.
    d) If new row inserted, publish to Kafka. If Kafka publish succeeds => return 201 accepted. If publish fails => update DB row with kafka_status='pending' and return 202.
  • Add configuration: REDIS_URL, IDEMPOTENCY_CACHE_TTL env vars; update requirements.txt/pyproject to include redis and cachetools fallback.
  • Add a small LRU cache fallback implementation when Redis is not available (for single-instance development only). Document that Redis is required in production.
  1. Consumer changes
  • Locate the consumer that reads from the mpesa-transactions topic and writes to DB. Update its DB write to be idempotent using INSERT ... ON CONFLICT DO NOTHING, and ensure offsets are committed only after DB write.
  1. Tests
  • Add integration tests that POST the same webhook payload twice and assert the DB has only one row and the second response is 200 duplicate.
  • Add unit tests for the Redis/local-cache branch.
  1. Dependencies & configuration
  • Add redis to requirements.txt (or pyproject). Add cachetools as optional fallback.
  • Update README / deployment docs to mention Redis and the migration checklist for adding the unique index (checking duplicates and running CONCURRENTLY when safe).

Implementation notes and safety

  • The migration must not blindly create a unique index if duplicates exist; instead it should fail and print SQL to find duplicates. Example pre-check:
    SELECT trans_id, count() FROM transactions GROUP BY trans_id HAVING count() > 1 LIMIT 1;
  • If duplicates exist, the PR should include a helper script tools/dedupe_transid.py that can be run to dedupe or a SQL snippet to review duplicates and remove them safely.
  • Using Redis.set(..., nx=True) is an atomic guard to avoid concurrent double-insert attempts. The DB unique index is the ultimate safety net.
  • When Kafka publish fails, prefer marking the DB row kafka_status='pending' so a worker can re-publish rather than deleting the row and losing data.

Files to add or modify (suggested)

  • migrations/20260903_add_unique_index_transid.sql (new)
  • alembic/versions/_add_unique_index_transid.py (new)
  • tools/dedupe_transid.py (new) — optional helper to find/remove duplicates
  • services/webhooks/handler.py (modify) — implement cache check + DB upsert + kafka publish
  • services/consumer/transactions_consumer.py (modify) — implement idempotent insert
  • requirements.txt or pyproject.toml (modify) — add redis and cachetools
  • tests/integration/test_webhook_idempotency.py (new)
  • README.md (modify) — docs for migration and Redis requirement

Code examples (to guide exact behavior)

SQL migration example (file: migrations/20260903_add_unique_index_transid.sql):

-- migrations/20260903_add_unique_index_transid.sql
-- Pre-check for duplicates. If this returns any rows, stop and dedupe before running.
SELECT trans_id, COUNT() as cnt FROM transactions GROUP BY trans_id HAVING COUNT() > 1 LIMIT 1;

-- If no duplicates exist, create a unique index concurrently (Postgres). Note: CONCURRENTLY cannot be run inside a transaction block.
-- Run the following on the DB directly or ensure your migration runner executes it without an enclosing transaction:
CREATE UNIQUE INDEX CONCURRENTLY IF NOT EXISTS ux_transactions_trans_id ON transactions (trans_id);

Webhook handler example (Python/Flask + SQLAlchemy + Redis):

services/webhooks/handler.py

from redis import Redis
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy import func
from myapp.db import SessionLocal, Transaction
from myapp.kafka import producer
import json, os

REDIS_URL = os.getenv('REDIS_URL', 'redis://localhost:6379/0')
CACHE_TTL_SECONDS = int(os.getenv('IDEMPOTENCY_CACHE_TTL', '300'))
redis_client = Redis.from_url(REDIS_URL)

def handle_webhook(request_json):
trans_id = extract_trans_id(request_json)
if not trans_id:
return {'error': 'missing trans_id'}, 400

cache_key = f"mpesa:trans:{trans_id}"
ok = False
try:
    ok = redis_client.set(cache_key, '1', ex=CACHE_TTL_SECONDS, nx=True)
except Exception:
    # redis unavailable -> fallback to local TTL cache or continue: prefer safety but allow best-effort
    ok = try_set_local_cache(cache_key)

if not ok:
    return {'status': 'duplicate'}, 200

db = SessionLocal()
try:
    insert_stmt = pg_insert(Transaction).values(
        trans_id=trans_id,
        amount=request_json.get('amount'),
        raw_payload=json.dumps(request_json),
        created_at=func.now()
    ).on_conflict_do_nothing(index_elements=['trans_id'])

    res = db.execute(insert_stmt)
    db.commit()
    if getattr(res, 'rowcount', 0) == 0:
        return {'status': 'duplicate'}, 200

    # Publish to Kafka
    producer.produce(topic='mpesa-transactions', key=str(trans_id), value=json.dumps(request_json).encode('utf-8'))
    producer.flush(timeout=5)
    return {'status': 'accepted'}, 201

except Exception as ex:
    db.rollback()
    # mark DB row pending (if present) or log
    db.query(Transaction).filter(Transaction.trans_id == trans_id).update({'kafka_status': 'pending'})
    db.commit()
    return {'status': 'queued_for_retry'}, 202
finally:
    db.close()

Consumer example (idempotent write):

services/consumer/transactions_consumer.py

from sqlalchemy.dialects.postgresql import insert as pg_insert
from myapp.db import SessionLocal, Transaction
from sqlalchemy import func

def process_message(payload):
trans_id = payload['trans_id']
db = SessionLocal()
try:
stmt = pg_insert(Transaction).values(
trans_id=trans_id,
amount=payload.get('amount'),
raw_payload=payload,
created_at=func.now()
).on_conflict_do_nothing(index_elements=['trans_id'])
res = db.execute(stmt)
db.commit()
if getattr(res, 'rowcount', 0) == 0:
logger.info('Duplicate message ignored for %s', trans_id)
finally:
db.close()

Testing plan

  • tests/integration/test_webhook_idempotency.py uses the running app and database fixture. It POSTs identical payload twice and asserts:
    • first response is 201 accepted
    • second response is 200 duplicate
    • DB has exactly one row for trans_id

PR and review

Additional notes for the implementer

  • The coding agent should search the repository for the existing webhook handler and consumer files and apply modifications in-place where appropriate.
  • If multiple webhook handlers exist, update the ones used in production (identify the file referenced by the web server entrypoint). If uncertain, update the files that contain the string "mpesa" or "webhook" and add TODO comments to any ambiguous locations.
  • Update requirements.txt/pyproject to include redis and cachetools, and run tests.

Deliverables

  • Branch: feat/idempotency-transid-20260903
  • Files added/updated per the list above
  • Alembic migration and plain SQL migration
  • Integration tests
  • A Pull Request opened against main with the title: "Enforce idempotency: DB unique index, webhook cache check, idempotent consumer" and description summarizing the changes and referencing issue Webhook idempotency for Daraja C2B confirmation callback #1.

The following is the prior conversation context from the user's chat exploration (may be truncated):

User: can you check for any modifications

Assistant: Yes — I checked the repository for recent modifications. I fetched the repo metadata and the latest commits on the default branch (main). Summary and the most recent commits below.

What I did

  • Queried the repository metadata (owner: Victor-Kipruto-Rop, repo: Real_Time_Transaction_Streaming-MPESA-).
  • Listed the latest commits (first page, 10 results).

Quick findings

  • Repository default branch: main
  • Last push: 2026-08-28T12:01:27Z
  • Open issues: 3
  • Language: Python
  • There are more commits beyond the first page — see link below.

Latest commits (most recent first)

  • 9924c5e — 2026-08-28T12:01:02Z — Victor Kipruto — "Harden MPESA transaction streaming pipeline"
    9924c5e
  • 3549970 — 2026-07-11T21:21:56Z — Victor kipruto Rop — "Merge pull request Checking for errors in the codebase #3 from Victor-Kipruto-Rop/copilot/check-for-erros — Checking for errors in the codebase"
    3549970
  • c62bf36 — 2026-07-11T20:18:41Z — copilot-swe-agent[bot] — "Changes before error encountered" (agent logs: session bd31c3de-...)
    c62bf36
  • 88e81a5 — 2026-07-11T20:17:16Z — copilot-swe-agent[bot] — "fix: continue compatibility fixes for unit test expectations"
    88e81a5
  • 58b2715 — 2026-07-11T20:14:51Z — copilot-swe-agent[bot] — "fix: restore legacy ingestion compatibility APIs"
    58b2715
  • e6dc017 — 2026-07-11T20:10:39Z — copilot-swe-agent[bot] — "fix: clean lint violations and restore missing compatibility symbols"
    e6dc017
  • c2bb864 — 2026-07-11T20:09:06Z — copilot-swe-agent[bot] — "fix: resolve lint issues and add flask dependency"
    c2bb864
  • ad329ff — 2026-07-11T19:40:32Z — Victor kipruto Rop — "Add files via upload"
    ad329ff
  • bd8e74e — 2026-07-11T19:38:45Z — Victor kipruto Rop — "Delete SECURITY.md"
    bd8e74e
  • 79944bd — 2026-06-15T05:05:01Z — Victor Kipruto — "Refactor configuration, improve data quality, organize tests, and update security policy"
    79944bd

More results and repository view

Would you like me to:

  • Show the files changed / full diff for the latest commit (9924c5e)?
  • Fetch the next page of commits (older commits)?
  • List open PRs or show changes for a specific PR/commit?

User: Critical Issues (Security & Reliability)
Webhook Idempotency (Issue #1) — OPEN

Duplicate TransID deliveries can create duplicate transaction rows
Missing unique DB constraint on TransID
No idempotency check before Kafka publish
Fix needed: Add DB-level unique index + in-memory cache check + consumer-side dedup

Assistant: Start with short recommendation, then concrete, copy-paste changes you can apply (migration + webhook handler + consumer change). These enforce DB uniqueness, add an atomic cache check before publish, and make the consumer idempotent.

Why: a DB unique index prevents persistent duplicates; an atomic cache check (Redis SET NX) short-circuits duplicate webhook deliveries before Kafka publish; consumer-side INSERT ... ON CONFLICT DO NOTHING guarantees duplicate messages that slip through won't create duplicate rows.

  1. DB migration — add unique index on TransID
  • File (new): migrations/20260903_add_unique_index_transid.sql
-- Add a unique index on transactions.trans_id to prevent duplicate rows
-- Replace "transactions" and "trans_id" with real table/column names if different.

BEGIN;

-- If you use PostgreSQL and the column is nullable, you may want a partial index that only applies to non-null trans_id:
CREATE UNIQUE INDEX IF NOT EXISTS ux_transactions_trans_id
  ON transactions (trans_id);

-- If you prefer to make the column NOT NULL + unique, do that in a separate migration:
-- ALTER TABLE transactions ALTER COLUMN trans_id SET NOT NULL;

COMMIT;

If you use Alembic, here's an example revision you can paste into alembic/versions/:

"""Add unique index on transactions.trans_id"""

from alembic import op

# revision identifiers, used by Alembic.
revision = "<rev>"
down_revision = "<prev_rev>"
branch_labels = None
depends_on = None

def upgrade():
    # Adjust table and column names to match your schema
    op.create_index("ux_transactions_trans_id", "transactions", ["trans_id"], unique=True)

def downgrade():
    op.drop_index("ux_transactions_trans_id", table_name="transactions")
  1. Webhook handler — check cache + attempt atomic DB insert (or check) before publishing to Kafka
  • File (edit): services/webhooks/handler.py (or wherever webhook processing lives)

Replace/augment the handler that currently:

  • accepts webhook -> publishes to Kafka -> inserts into DB (or inserts after publish)

with the following flow:

  1. Normalize/extract trans_id from the webhook payload.
  2. Do an atomic cache check: redis.set(key, 1, ex=TTL, nx=True) — if it returns False, treat as duplicate and respond 200.
  3. Attempt to insert to DB with INSERT ... ON CONFLICT DO NOTHING (or check DB existence). If DB insert fails due to unique constraint, treat as duplicate.
  4. If insert succeeded (new row), publish to Kafka.
  5. If Kafka publish fails, roll back DB insert or mark transaction for retry; ensure cache TTL prevents reprocessing during retries.

Example code (SQLAlchemy + Redis + confluent-kafka). Adjust imports/objects to fit your repo:

from redis import Redis
import os
from sqlalchemy.dialects.postgresql import insert as pg_insert
from sqlalchemy.exc import IntegrityError
from myapp.db import SessionLocal, Transaction  # adjust to your codebase
from myapp.kafka import producer  # adjust to your kafka wrapper
from flask import request, jsonify

REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
CACHE_TTL_SECONDS = int(os.getenv("IDEMPOTENCY_CACHE_TTL", "300"))  # 5 minutes

redis_client = Redis.from_url(REDIS_URL)

def handle_webhook():
    payload = request.json
    trans_id = payload.get("trans_id") or payload.get("TransactionId")  # normalize field name
    if not trans_id:
        return jsonify({"error": "missing trans_id"}), 400

    cache_key = f"mpesa:trans:{trans_id}"

    # Atomic cache check: set NX (only set if not exists)
    was_set = redis_client.set(cache_key, "1", ex=CACHE_TTL_SECONDS, nx=True)
    if not was_set:
        # We already processed this trans_id recently: safe idempotent response
        return jsonify({"status": "duplicate"}), 200

    db = SessionLocal()
    try:
        # Try an atomic insert which relies on DB unique constraint.
        # This assumes Transaction model maps to "transactions" and has trans_id column.
        insert_stmt = pg_insert(Transaction).values(
            trans_id=trans_id,
            amount=payload.get("amount"),
            raw_payload=payload,
            created_at=func.now()
        ).on_conflict_do_nothing(index_elements=["trans_id"])

        result = db.execute(insert_stmt)
        db.commit()

        # If no row inserted (result.rowcount == 0), it's a duplicate in DB
        if getattr(result, "rowcount", None) == 0:
            # duplicate — safe to return
            return jsonify({"status": "duplicate"}), 200

        # New row inserted; now publish to Kafka
        try:
            producer.produce(
                topic="mpesa-transactions",
                key=str(trans_id),
                value=json.dumps(payload).encode("utf-8")
            )
            producer.flush(timeout=5)
        except Exception:
            # Kafka publish failed — cleanup or mark for retry.
            # Option A: delete the DB row so other deliveries may re-insert (risky if consumer relies on DB).
            # Option B: mark the DB row status='pending_kafka' so a background job retries.
            # I'll mark for retry:
            db.query(Transaction).filter(Transaction.trans_id == trans_id).update({"kafka_status": "pending"})
            db.commit()
            return jsonify({"status": "queued_for_retry"}), 202

        return jsonify({"status": "accepted"}), 201

    except IntegrityError:
        db.rollback()
        # Duplicate at DB level — idempotent response
        return jsonify({"status": "duplicate"}), 200
    finally:
        db.close()

Notes:

  • Using Redis.set(..., nx=True) is atomic and prevents a race where two webhooks attempt insertion concurrently.
  • The DB unique index is the ultimate guard—even with race conditions this prevents duplicates.
  • When Kafka publish fails, prefer marking DB row for retry rather than deleting the row (so you don’t lose data); a background worker can re-publish records with kafka_status='pending'.
  1. In-memory fallback (if you don't have Redis) — small LRU cache for short window
  • File (new/optional): services/webhooks/cache.py
# Simple LRU fallback when Redis isn't available. Not recommended for multi-instance deployments.
from cachetools import TTLCache

local_cache = TTLCache(maxsize=10000, ttl=300)  # fits recent trans_ids

def try_set_local_cache(key):
    if key in local_cache:
        return False
    local_cache[key] = True
    return True

Then use try_set_local_cache(cache_key) like the redis check. But use Redis in production (multi-instance safe).

  1. Consumer-side dedup — idempotent insert (safe if duplicate messages are consumed)
  • File (edit): services/consumer/transactions_consumer.py

Replace the consumer DB write with an idempotent insert:

from sqlalchemy.dialects.postgresql import insert as pg_insert
from myapp.db import SessionLocal, Transaction
import json

def process_message(message):
    payload = json.loads(message.value())
    trans_id = payload["trans_id"]
    db = SessionLocal()
    try:
        stmt = pg_insert(Transaction).values(
            trans_id=trans_id,
            amount=payload.get("amount"),
            raw_payload=payload,
            created_at=func.now()
        ).on_conflict_do_nothing(index_elements=["trans_id"])
        result = db.execute(stmt)
        db.commit()
        if getattr(result, "rowcount", None) == 0:
            # Duplicate, optionally log and ack
            logger.info("Duplicate message ignored for trans_id=%s", trans_id)
        else:
            logger.info("Inserted transaction %s", trans_id)
    finally:
        db.close()

If your consumer commits offsets only after successful DB write, duplicates will be rare, but the ON CONFLICT guarantees safety.

  1. Tests
  • Add integration tests: send duplicate webhook payloads, assert only one DB row is present and the second call returns 200 duplicate.
  • Add unit tests for Redis fallback and for consumer INSERT ON CONFLICT behavior.

Risk considerations & deployment

  • Adding a unique index can fail if there are existing duplicate rows. Before running migration:
    • Find duplicates: SELECT trans_id, count() FROM transactions GROUP BY trans_id HAVING count() > 1;
    • Clean duplicates (remove or dedupe) or create a unique index concurrently after removing duplicates:
      • CREATE UNIQUE INDEX CONCURRENTLY ux_transactions_trans_id ON transactions (trans_id);
      • Note: CONCURRENTLY requires no other transaction making schema changes and cannot run inside a transaction block.
  • If you make trans_id NOT NULL, verify all existing rows have non-null values.
  • Use small cache TTL (e.g., 5 minutes) to avoid long-lived keys that block legitimate retries.

What I did here

  • Drafted a migration (SQL + Alembic sample).
  • Wrote copy-paste-ready handler and consumer code snippets that implement the cache check + DB unique insertion + consumer dedupe.
  • Suggested a Redis-based approach (recommended) and a local LRU fallback.

Want me to…

  1. Implement these changes across the repository (add migration + modify webhook handler + modify consumer) and open a pull request for review.
  2. Create the DB migration only (Alembic revision or SQL file) and a checklist for safely applying it to production (duplicate cleanup + CONCURRENTLY guidance).
  3. Add automated integration tests that POST duplicate webhooks and assert single DB row + correct kafka behavior.

Tell me which option (1-3) you want and I will proceed.

User: Implement these changes across the repository (add migration + modify webhook handler + modify consumer) and open a pull request for review.

Assistant: I'm waiting for your confirmation to proceed with implementing the changes and opening the pull request—please accept the prompt so I can continue.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Webhook idempotency for Daraja C2B confirmation callback

2 participants