winx_prinx-api / server /cache.js
Sasha
Fix: only cache successful 2xx responses and validate streams array structure
c22c35b
Raw
History Blame
9.83 kB
/**
* production-ready In-Memory Caching and Rate-Limiting Subsystem
* Designed to protect the Supabase DB and Express server from spikes in traffic
*/
import dotenv from 'dotenv';
dotenv.config();
// Configuration flags
const IS_CACHE_ENABLED = process.env.CACHE_ENABLED !== 'false';
const DEFAULT_TTL_MS = (parseInt(process.env.CACHE_DEFAULT_TTL, 10) || 15) * 1000;
const CACHE_MAX_KEYS = parseInt(process.env.CACHE_MAX_KEYS, 10) || 1000;
const RATE_LIMIT_MAX = parseInt(process.env.RATE_LIMIT_MAX_REQUESTS, 10) || 60;
const RATE_LIMIT_WINDOW_MS = parseInt(process.env.RATE_LIMIT_WINDOW_MS, 10) || 60000;
// Metrics tracker
const metrics = {
hits: 0,
misses: 0,
setCount: 0,
blockedRequests: 0,
};
// =========================================================================
// CACHE MANAGER CLASS
// =========================================================================
class CacheManager {
constructor() {
this.store = new Map();
// Background garbage collection sweep every 30 seconds
this.gcTimer = setInterval(() => {
this.evictExpired();
}, 30000).unref(); // .unref() allows Node process to exit cleanly if running tests
}
set(key, value, ttlMs = DEFAULT_TTL_MS) {
if (!IS_CACHE_ENABLED) return;
const now = Date.now();
const expiresAt = now + ttlMs;
if (this.store.has(key)) {
// Delete existing to update its insertion order position (LRU)
this.store.delete(key);
} else if (this.store.size >= CACHE_MAX_KEYS) {
// Evict the least recently used (oldest) key
const oldestKey = this.store.keys().next().value;
if (oldestKey !== undefined) {
this.store.delete(oldestKey);
console.log(`[CacheManager LRU] Evicted oldest key: ${oldestKey}`);
}
}
this.store.set(key, {
value,
expiresAt,
});
metrics.setCount++;
}
get(key) {
if (!IS_CACHE_ENABLED) return null;
const entry = this.store.get(key);
if (!entry) return null;
// Refresh insertion order position for LRU
this.store.delete(key);
this.store.set(key, entry);
// Return the entry even if expired. Freshness logic is handled by middleware
return entry;
}
del(key) {
return this.store.delete(key);
}
flush() {
this.store.clear();
console.log('[CacheManager] Cache flushed successfully.');
}
evictExpired() {
const now = Date.now();
let evictedCount = 0;
for (const [key, entry] of this.store.entries()) {
// Keep stale data in memory for up to 24 hours so it can be served during heavy compute spikes
if (now > entry.expiresAt + 24 * 60 * 60 * 1000) {
this.store.delete(key);
evictedCount++;
}
}
if (evictedCount > 0) {
console.log(`[CacheManager GC] Swept and evicted ${evictedCount} expired cache keys.`);
}
}
getStats() {
const totalKeys = this.store.size;
const totalRequests = metrics.hits + metrics.misses;
const hitRatio = totalRequests > 0 ? ((metrics.hits / totalRequests) * 100).toFixed(1) + '%' : '0%';
return {
enabled: IS_CACHE_ENABLED,
totalKeys,
hits: metrics.hits,
misses: metrics.misses,
hitRatio,
setCount: metrics.setCount,
blockedRequests: metrics.blockedRequests,
};
}
}
export const cache = new CacheManager();
// =========================================================================
// RATE LIMITER CLASS
// =========================================================================
class RateLimiter {
constructor() {
this.clients = new Map();
// Background sweeper to clear expired rate-limit buckets every 30 seconds
this.gcTimer = setInterval(() => {
this.evictExpired();
}, 30000).unref();
}
/**
* Check if client IP is within limits
* @returns {object} { allowed: boolean, remaining: number, limit: number, resetTime: number }
*/
check(ip) {
const now = Date.now();
let client = this.clients.get(ip);
if (!client || now > client.resetTime) {
// Create new window
const resetTime = now + RATE_LIMIT_WINDOW_MS;
client = {
hits: 1,
resetTime,
};
this.clients.set(ip, client);
return {
allowed: true,
remaining: RATE_LIMIT_MAX - 1,
limit: RATE_LIMIT_MAX,
resetTime,
};
}
client.hits++;
const remaining = Math.max(0, RATE_LIMIT_MAX - client.hits);
const allowed = client.hits <= RATE_LIMIT_MAX;
return {
allowed,
remaining,
limit: RATE_LIMIT_MAX,
resetTime: client.resetTime,
};
}
evictExpired() {
const now = Date.now();
let count = 0;
for (const [ip, data] of this.clients.entries()) {
if (now > data.resetTime) {
this.clients.delete(ip);
count++;
}
}
if (count > 0) {
console.log(`[RateLimiter GC] Evicted ${count} expired rate-limit records.`);
}
}
}
const rateLimiter = new RateLimiter();
// =========================================================================
// EXPRESS MIDDLEWARE
// =========================================================================
const pendingRequests = new Map();
/**
* Express Middleware protecting route with Rate-Limiting, Coalescing, and Stale-While-Revalidate Cache
* @param {number} ttlSeconds Custom TTL override in seconds
*/
export function cacheMiddleware(ttlSeconds) {
const ttlMs = ttlSeconds ? ttlSeconds * 1000 : DEFAULT_TTL_MS;
return async (req, res, next) => {
// 1. Bypass check if cache is disabled globally
if (!IS_CACHE_ENABLED) {
return next();
}
// Extract real client IP (supports Express trust proxy)
const ip = req.ip || req.headers['x-forwarded-for'] || req.socket.remoteAddress;
// 2. Perform Rate Limiting
const limitStatus = rateLimiter.check(ip);
// Set standard rate limit headers
res.setHeader('X-RateLimit-Limit', limitStatus.limit);
res.setHeader('X-RateLimit-Remaining', limitStatus.remaining);
res.setHeader('X-RateLimit-Reset', Math.ceil(limitStatus.resetTime / 1000));
if (!limitStatus.allowed) {
metrics.blockedRequests++;
const retryAfter = Math.ceil((limitStatus.resetTime - Date.now()) / 1000);
res.setHeader('Retry-After', retryAfter);
console.warn(`[RateLimiter] Blocked IP: ${ip} for too many requests on path: ${req.path}`);
return res.status(429).json({
error: 'Too Many Requests',
message: 'Превышен лимит запросов к серверу. Пожалуйста, подождите.',
retryAfter,
});
}
// 3. Perform Cache Lookup
const sortedQueries = Object.keys(req.query)
.sort()
.map(k => `${k}=${req.query[k]}`)
.join('&');
const cacheKey = sortedQueries ? `${req.path}?${sortedQueries}` : req.path;
const entry = cache.get(cacheKey);
const now = Date.now();
// CASE A: Cache is FRESH
if (entry && now <= entry.expiresAt) {
metrics.hits++;
const remainingTtlSecs = Math.max(0, Math.ceil((entry.expiresAt - now) / 1000));
res.setHeader('X-Cache', 'HIT');
res.setHeader('X-Cache-TTL-Remaining', remainingTtlSecs);
return res.json(entry.value);
}
// CASE B: Cache is STALE or MISS, and another request is already computing it
if (pendingRequests.has(cacheKey)) {
if (entry) {
// Stale-While-Revalidate: Return old data immediately to avoid waiting
metrics.hits++;
res.setHeader('X-Cache', 'STALE');
return res.json(entry.value);
} else {
// Request Coalescing: No old data exists, so wait for the pending request to finish
metrics.misses++;
res.setHeader('X-Cache', 'COALESCED');
try {
const value = await pendingRequests.get(cacheKey);
return res.json(value);
} catch (e) {
return res.status(500).json({ error: 'Internal Server Error computing cache' });
}
}
}
// CASE C: Cache is STALE or MISS, and WE are the first request to trigger compute
metrics.misses++;
res.setHeader('X-Cache', 'MISS');
// Create a pending Promise that others can wait on
let resolvePending;
let rejectPending;
const pendingPromise = new Promise((resolve, reject) => {
resolvePending = resolve;
rejectPending = reject;
});
pendingRequests.set(cacheKey, pendingPromise);
// Timeout safety net to prevent frozen promises if next() hangs forever
const safetyTimeout = setTimeout(() => {
if (pendingRequests.get(cacheKey) === pendingPromise) {
pendingRequests.delete(cacheKey);
rejectPending(new Error('Cache compute timeout'));
}
}, 5 * 60 * 1000);
// Intercept res.json to store the result
const originalJson = res.json;
res.json = function (body) {
clearTimeout(safetyTimeout);
res.json = originalJson; // Restore
// Save fresh data only on successful 2xx responses and check array structure for streams
if (res.statusCode >= 200 && res.statusCode < 300) {
let shouldCache = true;
// If this is the streams list endpoint, ensure the body is a valid array
if (req.path === '/api/streams' && !Array.isArray(body)) {
shouldCache = false;
}
if (shouldCache) {
cache.set(cacheKey, body, ttlMs);
}
}
// Resolve the promise to unblock anyone who was coalesced (waiting)
resolvePending(body);
if (pendingRequests.get(cacheKey) === pendingPromise) {
pendingRequests.delete(cacheKey);
}
return res.json(body);
};
next();
};
}