DEV Community

subashthiruppathy
subashthiruppathy

Posted on

Implementing Webhook Idempotency with Node.js, Redis, and PostgreSQL

Implementing Webhook Idempotency with Node.js, Redis, and PostgreSQL

In financial and transaction-heavy backends, external webhook deliveries (from providers like Stripe, Razorpay, or banking partners) must be handled with extreme care.

Network timeouts, server retries, and transit jitter mean that third-party gateways will frequently send the exact same event payload multiple times. If your webhook handler isn't strictly idempotent, duplicate requests will lead to corrupted state, duplicate invoice generation, or crediting user balances twice.

In this guide, we'll design a production-grade webhook ingestion handler using Node.js (Express), Redis for distributed concurrency locking, and PostgreSQL for transactional persistence.


The Two Failure Modes of Naive Webhook Handlers

  1. Sequential Duplicates: The payment gateway retries a webhook 30 seconds later because your initial response timed out. A naive database read-then-write pattern (SELECT then INSERT) often misses the window if the worker is still running.
  2. Concurrent Race Conditions: Two identical delivery attempts hit your cluster at the exact same millisecond across different server instances. Both instances query the database, find no existing record, and proceed to execute the transaction twice.

To solve both issues, we need a two-layer defense:

  • In-Memory Concurrency Lock (Redis): Intercepts concurrent executions before they touch business logic.
  • Persistent Unique Constraint (PostgreSQL): Guarantees zero duplicate ledger records even if Redis restarts.

Step 1: Preserving Raw Buffers for Cryptographic Verification

Before parsing JSON, you must verify the signature header (HMAC-SHA256) sent by the provider. If you verify against parsed JSON, slight formatting differences will cause hash mismatches.

import express from 'express';
import crypto from 'crypto';

const app = express();

// Preserve the raw request buffer for signature checks
app.use(express.json({
  verify: (req, res, buf) => {
    req.rawBody = buf;
  }
}));

function verifySignature(req, webhookSecret) {
  const signature = req.headers['x-signature'];
  const expectedSignature = crypto
    .createHmac('sha256', webhookSecret)
    .update(req.rawBody)
    .digest('hex');

  // Use timingSafeEqual to protect against timing attacks
  return crypto.timingSafeEqual(
    Buffer.from(signature || '', 'hex'),
    Buffer.from(expectedSignature, 'hex')
  );
}
Enter fullscreen mode Exit fullscreen mode

Step 2: Concurrency Control Using Redis Atomic Locks

If two identical delivery attempts reach separate server instances simultaneously, traditional read-then-write checks (SELECT before INSERT) will encounter race conditions.

Using Redis SET with NX (Not Exists) and PX (TTL in milliseconds) provides an atomic lock:

import Redis from 'ioredis';
const redis = new Redis(process.env.REDIS_URL);

async function acquireLock(eventId, ttlMs = 10000) {
  const lockKey = `lock:webhook:${eventId}`;
  // Returns 'OK' if key was set, null if it already exists
  const result = await redis.set(lockKey, 'locked', 'NX', 'PX', ttlMs);
  return result === 'OK';
}

async function releaseLock(eventId) {
  await redis.del(`lock:webhook:${eventId}`);
}
Enter fullscreen mode Exit fullscreen mode

Step 3: PostgreSQL Schema and Atomic Ledger Handling

The persistent layer requires a unique constraint on the provider's event_id and tracks lifecycle statuses (processing, completed, failed):

CREATE TABLE webhook_events (
    id SERIAL PRIMARY KEY,
    event_id VARCHAR(128) NOT NULL UNIQUE,
    event_type VARCHAR(64) NOT NULL,
    status VARCHAR(32) NOT NULL DEFAULT 'processing',
    payload JSONB NOT NULL,
    created_at TIMESTAMP WITH TIME ZONE DEFAULT CURRENT_TIMESTAMP
);

CREATE TABLE user_balances (
    user_id UUID PRIMARY KEY,
    balance_cents BIGINT NOT NULL DEFAULT 0
);
Enter fullscreen mode Exit fullscreen mode

Idempotent Transaction Handler

import pg from 'pg';
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL });

async function processWebhookEvent(event) {
  const client = await pool.connect();

  try {
    await client.query('BEGIN');

    // Attempt insert; ON CONFLICT DO NOTHING guarantees safe deduplication
    const insertQuery = `
      INSERT INTO webhook_events (event_id, event_type, payload, status)
      VALUES ($1, $2, $3, 'processing')
      ON CONFLICT (event_id) DO NOTHING
      RETURNING id;
    `;
    const res = await client.query(insertQuery, [
      event.id,
      event.type,
      JSON.stringify(event.data)
    ]);

    // If rowCount === 0, event was already inserted by a prior delivery
    if (res.rowCount === 0) {
      await client.query('ROLLBACK');
      return { status: 'already_processed' };
    }

    // Process business logic atomically inside transaction
    if (event.type === 'payment.succeeded') {
      const { userId, amountCents } = event.data;
      await client.query(
        `UPDATE user_balances
         SET balance_cents = balance_cents + $1
         WHERE user_id = $2`,
        [amountCents, userId]
      );
    }

    // Mark event status as completed
    await client.query(
      `UPDATE webhook_events SET status = 'completed' WHERE event_id = $1`,
      [event.id]
    );

    await client.query('COMMIT');
    return { status: 'success' };
  } catch (err) {
    await client.query('ROLLBACK');
    throw err;
  } finally {
    client.release();
  }
}
Enter fullscreen mode Exit fullscreen mode

Step 4: Putting It Together: The Express Controller

app.post('/api/webhooks', async (req, res) => {
  const event = req.body;

  // 1. Authenticate origin
  if (!verifySignature(req, process.env.WEBHOOK_SECRET)) {
    return res.status(401).send('Invalid signature');
  }

  // 2. Concurrency check
  const hasLock = await acquireLock(event.id);
  if (!hasLock) {
    return res.status(200).json({ message: 'Event currently being processed' });
  }

  try {
    // 3. Execute atomic transaction
    const result = await processWebhookEvent(event);
    return res.status(200).json(result);
  } catch (error) {
    console.error(`Processing error for event ${event.id}:`, error);
    return res.status(500).send('Internal Server Error');
  } finally {
    await releaseLock(event.id);
  }
});
Enter fullscreen mode Exit fullscreen mode

Top comments (0)