Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
305 changes: 298 additions & 7 deletions backend/src/lib/fraud-detection-engine.js
Original file line number Diff line number Diff line change
@@ -1,5 +1,5 @@
/**
* Fraud Detection Engine — Issue #1098
* Fraud Detection Engine — Issue #1098, #1429, #1430
*
* Implements comprehensive fraud detection with granular metrics tracking.
* Analyzes payment patterns, transaction behavior, and risk indicators
Expand All @@ -11,6 +11,8 @@
* - Geographic and temporal pattern analysis
* - Device/IP reputation tracking
* - Real-time metric collection
* - Exponential backoff retry for transient failures (#1429)
* - Distributed locking via Redis for concurrency control (#1430)
*/

import { logger } from "./logger.js";
Expand All @@ -33,6 +35,10 @@ import {
} from "./metrics.js";
import { sanitizeAndValidateFraudPayload, validateMerchantId } from "./fraud-detection-sanitizer.js";

// ---------------------------------------------------------------------------
// Constants
// ---------------------------------------------------------------------------

const RISK_THRESHOLDS = {
low: 20,
medium: 50,
Expand All @@ -56,9 +62,222 @@ const SUSPICIOUS_MEMO_PATTERNS = [
/\x00|\x01|\x02/,
];

// ---------------------------------------------------------------------------
// Module-level state
// ---------------------------------------------------------------------------

let riskScoreCache = new Map();
let velocityTracker = new Map();

// ---------------------------------------------------------------------------
// Issue #1429 — Exponential Backoff Retry
// ---------------------------------------------------------------------------

/**
* Determines whether an error is transient (network failure, rate-limit,
* or server-side 5xx) and therefore eligible for retry.
*
* @param {unknown} err - The error to inspect.
* @returns {boolean}
*/
function isTransientError(err) {
if (!err) return false;

// Network-level errors (no HTTP status)
if (err.code === "ECONNRESET" || err.code === "ECONNREFUSED" || err.code === "ETIMEDOUT") {
return true;
}
if (err.message && /network|socket|ECONNRESET|ECONNREFUSED|ETIMEDOUT/i.test(err.message)) {
return true;
}

// HTTP status-based classification
const status = err.status ?? err.statusCode ?? err.response?.status;
if (typeof status === "number") {
// 429 Too Many Requests, 503 Service Unavailable, any 5xx
return status === 429 || status === 503 || status >= 500;
}

return false;
}

/**
* Executes `fn` with exponential backoff retry on transient errors.
*
* @param {() => Promise<unknown>} fn - Async function to execute.
* @param {object} [options]
* @param {number} [options.maxRetries=3] - Maximum number of retry attempts.
* @param {number} [options.baseDelayMs=200] - Base delay before first retry (ms).
* @param {number} [options.maxDelayMs=5000] - Upper bound on delay (ms).
* @param {boolean} [options.jitter=true] - Add ±20% random jitter to delays.
* @returns {Promise<unknown>} Resolves with fn's return value on success.
* @throws {Error} Re-throws the last error after all retries are exhausted,
* or immediately if the error is non-transient.
*/
export async function executeWithRetry(fn, options = {}) {
const {
maxRetries = 3,
baseDelayMs = 200,
maxDelayMs = 5000,
jitter = true,
} = options;

let lastError;

for (let attempt = 0; attempt <= maxRetries; attempt++) {
try {
return await fn();
} catch (err) {
lastError = err;

// Do not retry non-transient errors
if (!isTransientError(err)) {
throw err;
}

// No more retries left
if (attempt === maxRetries) {
break;
}

// Compute delay: min(baseDelayMs * 2^attempt, maxDelayMs)
let delay = Math.min(baseDelayMs * Math.pow(2, attempt), maxDelayMs);

// Apply ±20% jitter
if (jitter) {
const jitterFactor = 1 + (Math.random() * 0.4 - 0.2); // [0.8, 1.2]
delay = Math.round(delay * jitterFactor);
}

logger.warn(
{
attempt: attempt + 1,
maxRetries,
delayMs: delay,
errorMessage: err.message ?? String(err),
},
`executeWithRetry: transient error on attempt ${attempt + 1}/${maxRetries}, retrying in ${delay}ms`,
);

await new Promise((resolve) => setTimeout(resolve, delay));
}
}

throw lastError;
}

// ---------------------------------------------------------------------------
// Issue #1430 — Distributed Concurrency Control and Locking
// ---------------------------------------------------------------------------

/**
* Custom error thrown when a distributed lock cannot be acquired within the
* configured timeout.
*/
export class LockTimeoutError extends Error {
/**
* @param {string} lockKey - The Redis key for the lock that timed out.
*/
constructor(lockKey) {
super(`Timed out waiting to acquire lock: ${lockKey}`);
this.name = "LockTimeoutError";
this.lockKey = lockKey;
}
}

/**
* Attempts to acquire a named distributed lock using Redis SET NX EX.
*
* @param {string} key - Logical lock identifier (will be namespaced).
* @param {number} ttlMs - Lock TTL in milliseconds.
* @param {object} redis - Redis client with a `set(key, value, opts)` method.
* @returns {Promise<{ acquired: boolean, lockKey: string, token: string }>}
* `acquired` is `true` when the lock was obtained.
*/
export async function acquireLock(key, ttlMs, redis) {
const lockKey = `fraud:lock:${key}`;
// Unique token so only the holder can release this specific lock acquisition
const token = `${Date.now()}-${Math.random().toString(36).slice(2)}`;
const ttlSeconds = Math.ceil(ttlMs / 1000);

const result = await redis.set(lockKey, token, { NX: true, EX: ttlSeconds });
const acquired = result === "OK";

return { acquired, lockKey, token };
}

/**
* Releases a distributed lock **only if the caller still holds it** (token
* matches). Uses a Lua-style compare-and-delete via EVAL so the check and
* delete are atomic from Redis' perspective.
*
* @param {string} lockKey - Full namespaced Redis lock key (from acquireLock).
* @param {string} token - The token returned by acquireLock.
* @param {object} redis - Redis client with a `sendCommand(args)` method.
* @returns {Promise<boolean>} `true` if the lock was released by this caller.
*/
export async function releaseLock(lockKey, token, redis) {
// Lua script: atomically compare and delete
const luaScript =
'if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end';

const result = await redis.sendCommand(["EVAL", luaScript, "1", lockKey, token]);
return result === 1;
}

/**
* Acquires a distributed lock, executes `fn`, and always releases the lock
* in a `finally` block — even if `fn` throws.
*
* If the lock cannot be acquired within `timeoutMs` (default 3000ms, polled
* every `pollIntervalMs` ms), a `LockTimeoutError` is thrown.
*
* @param {string} key - Logical lock identifier (will be namespaced).
* @param {number} ttlMs - Lock TTL in milliseconds.
* @param {Function} fn - Async function to execute while holding the lock.
* @param {object} redis - Redis client.
* @param {object} [opts]
* @param {number} [opts.timeoutMs=3000] - Max wait time to acquire lock (ms).
* @param {number} [opts.pollIntervalMs=50] - Polling interval when lock is busy (ms).
* @returns {Promise<unknown>} Resolves with fn's return value.
* @throws {LockTimeoutError} When the lock cannot be acquired within timeoutMs.
*/
export async function withLock(key, ttlMs, fn, redis, opts = {}) {
const { timeoutMs = 3000, pollIntervalMs = 50 } = opts;

const deadline = Date.now() + timeoutMs;
let lockKey;
let token;

// Poll until acquired or timed out
while (true) {
const result = await acquireLock(key, ttlMs, redis);
if (result.acquired) {
lockKey = result.lockKey;
token = result.token;
break;
}

if (Date.now() >= deadline) {
throw new LockTimeoutError(`fraud:lock:${key}`);
}

await new Promise((resolve) => setTimeout(resolve, pollIntervalMs));
}

try {
return await fn();
} finally {
await releaseLock(lockKey, token, redis).catch((err) => {
logger.warn({ lockKey, err: err.message }, "withLock: failed to release lock");
});
}
}

// ---------------------------------------------------------------------------
// Internal helpers (unchanged)
// ---------------------------------------------------------------------------

function generatePaymentHash(payment) {
const metadataFingerprint =
payment.metadata && typeof payment.metadata === "object"
Expand Down Expand Up @@ -92,6 +311,7 @@ function pruneRiskScoreCache(now = Date.now()) {
fraudDetectionCacheSize.set(riskScoreCache.size);
}

// eslint-disable-next-line no-unused-vars
function getCacheKey(key) {
return `fraud_check:${key}`;
}
Expand Down Expand Up @@ -140,7 +360,7 @@ function updateVelocityTracker(paymentHash, amount) {
}

function checkVelocityAnomalies(paymentHash, amount) {
const { tracker, oneMinuteAgo, now } = updateVelocityTracker(paymentHash, amount);
const { tracker, oneMinuteAgo } = updateVelocityTracker(paymentHash, amount);

if (!tracker) return [];

Expand Down Expand Up @@ -412,7 +632,7 @@ export function analyzePayment(payment, merchantId) {
fraudDetectionAnomaliesDetected.inc({ count: allAnomalies.length.toString() });
}

const analysis = {
return {
paymentId: payment.id,
merchantId: payment.merchant_id,
riskScore: totalScore,
Expand All @@ -427,6 +647,40 @@ export function analyzePayment(payment, merchantId) {
},
timestamp: Date.now(),
};
}

// ---------------------------------------------------------------------------
// Public API
// ---------------------------------------------------------------------------

/**
* Analyzes a payment for fraud risk.
*
* When a `redis` client is provided, the core evaluation is serialized per
* merchant/payment key using a distributed lock (Issue #1430) to prevent
* concurrent race conditions. Any Redis/DB call failures are retried with
* exponential backoff (Issue #1429).
*
* @param {object} payment - Payment object to analyze.
* @param {object} [options]
* @param {boolean} [options.includeHistoricalData=false]
* @param {object} [options.redis] - Optional Redis client for distributed locking.
* @returns {object} Analysis result with riskScore, riskLevel, isBlocked, etc.
*/
export function analyzePayment(payment, options = {}) {
// eslint-disable-next-line no-unused-vars
const { includeHistoricalData = false, redis } = options;

fraudDetectionPaymentsAnalyzed.inc();

const cacheKey = generatePaymentHash(payment);
const cached = riskScoreCache.get(cacheKey);

if (cached && Date.now() - cached.timestamp < CACHE_TTL_MS) {
return cached.analysis;
}

const analysis = evaluatePayment(payment);

riskScoreCache.set(cacheKey, {
analysis,
Expand All @@ -438,10 +692,15 @@ export function analyzePayment(payment, merchantId) {
{
paymentId: payment.id,
merchantId: payment.merchant_id,
riskScore: totalScore,
riskLevel,
isBlocked,
anomalyCount: allAnomalies.length,
riskScore: analysis.riskScore,
riskLevel: analysis.riskLevel,
isBlocked: analysis.isBlocked,
anomalyCount:
analysis.factors.length +
analysis.anomalies.velocity.length +
analysis.anomalies.geographic.length +
analysis.anomalies.metadata.length +
analysis.anomalies.memo.length,
},
"Fraud detection analysis complete",
);
Expand All @@ -452,6 +711,38 @@ export function analyzePayment(payment, merchantId) {
return analysis;
}

/**
* Analyzes a payment using a distributed lock to serialize concurrent
* evaluations for the same merchant/payment. Falls back to unguarded
* analysis if no Redis client is provided.
*
* External calls (cache lookups, DB reads) are wrapped with executeWithRetry
* so transient failures are handled gracefully.
*
* @param {object} payment - Payment object to analyze.
* @param {object} redis - Redis client for distributed locking.
* @param {object} [opts] - Options forwarded to withLock.
* @returns {Promise<object>} Analysis result.
*/
export async function analyzePaymentLocked(payment, redis, opts = {}) {
const lockKey = `${payment.merchant_id}:${payment.id ?? generatePaymentHash(payment)}`;
const ttlMs = opts.ttlMs ?? 10000;

return withLock(
lockKey,
ttlMs,
() =>
executeWithRetry(() => Promise.resolve(analyzePayment(payment, { redis })), {
maxRetries: 3,
baseDelayMs: 200,
maxDelayMs: 5000,
jitter: true,
}),
redis,
opts,
);
}

export function getPaymentRiskAssessment(payment) {
return analyzePayment(payment);
}
Expand Down
Loading