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
-
Sequential Duplicates: The payment gateway retries a webhook 30 seconds later because your initial response timed out. A naive database read-then-write pattern (
SELECTthenINSERT) often misses the window if the worker is still running. - 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')
);
}
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}`);
}
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
);
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();
}
}
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);
}
});
Top comments (0)