Skip to content

[WIP] Add safe DB migration and tools for idempotent M-Pesa transactions - #6

Draft
Victor-Kipruto-Rop with Copilot wants to merge 1 commit into
mainfrom
copilot/featadd-unique-index-transid-20260903
Draft

[WIP] Add safe DB migration and tools for idempotent M-Pesa transactions#6
Victor-Kipruto-Rop with Copilot wants to merge 1 commit into
mainfrom
copilot/featadd-unique-index-transid-20260903

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.

This pull request adds a safe DB migration and supporting tools to enforce idempotency for M-Pesa transactions by introducing a unique database index on transactions.trans_id, plus helper scripts and documentation for safely applying the migration in production.

Changes included on branch feat/add-unique-index-transid-20260903:

  1. migrations/20260903_add_unique_index_transid.sql

    • SQL pre-check that lists duplicate trans_id groups (LIMIT 10). Includes recommended CREATE UNIQUE INDEX CONCURRENTLY command and fallback non-concurrent variant plus rollback command.
  2. alembic/versions/000000_add_unique_index_transid.py

    • Alembic revision example that performs a pre-check for duplicate trans_id values and aborts the migration if duplicates exist. Includes a pattern to issue CREATE UNIQUE INDEX CONCURRENTLY by committing around the statement with warnings in comments.
  3. tools/dedupe_transid.py

    • Batch dedupe helper that deletes older duplicate rows (keeps the newest by created_at), configurable via DATABASE_URL and DEDUPE_BATCH_SIZE. Intended to be run in a maintenance window after review.
  4. README.md (migration checklist & quick commands)

    • Documentation and step-by-step production checklist: backups, pre-checks, dedupe, running CREATE UNIQUE INDEX CONCURRENTLY, verification, and rollback notes.

Why

  • Duplicate webhook deliveries were creating duplicate transaction rows. Adding a unique index on trans_id prevents persistent duplicates. The pre-check and dedupe helper ensure the migration won't fail unexpectedly. Using CONCURRENTLY is recommended to avoid long locks on large tables.

Notes and next steps for reviewers

  • Confirm the target table name and column: the migration assumes transactions.trans_id. If your schema uses a different table/column, please request a change and I will update the migration and script.
  • This PR intentionally does not run application code changes (webhook handler or consumer) — it focuses on the DB migration and safe deployment steps. After the index is applied, follow-up PR(s) should implement cache-based idempotency and idempotent consumer inserts.
  • If you prefer the migration to be split into (a) pre-check Alembic migration and (b) manual CONCURRENTLY step, we can adjust the PR to mark the CONCURRENTLY creation as a manual step in deployment docs.

Testing

  • On staging: run the SELECT duplicate pre-check, run tools/dedupe_transid.py if duplicates exist, and then run the recommended CREATE UNIQUE INDEX CONCURRENTLY command. Verify inserts reject duplicates.

This PR fixes issue #1 (Webhook Idempotency) by adding the migration and guidance necessary to ensure DB-level uniqueness.

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.

User: @copilot Accepted Confirmation: Confirm cloud agent session

Assistant: Done — the Copilot coding agent has started implementing the idempotency changes you requested. You can track progress and see the work items at this task URL:

https://github.com/Victor-Kipruto-Rop/Real_Time_Transaction_Streaming-MPESA-/tasks/2c686932-61aa-4b46-8353-bc2733aa4ed1

What I started: the agent is working on a branch to add the DB migration (unique index + pre-check), modify the webhook handler to use an atomic cache (Redis SET NX) + idempotent DB insert before Kafka publish, make the consumer insert idempotent (INSERT ... ON CONFLICT DO NOTHING), add tests, and update dependencies and docs. The agent will open a pull request for review when the changes are ready and will update the task with PR links and progress.

I’ll post the PR link and a summary here as soon as the agent completes the work. If you want me to change the branch name, testing scope, or migration strategy (e.g., require CONCURRENTLY vs non-concurrent), tell me now and I’ll update the task.

User: Create the DB migration only (Alembic revision or SQL file) and a checklist for safely applying it to production (duplicate cleanup + CONCURRENTLY guidance).

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.

2 participants