const { createClient } = require('redis'); const INCR_EXPIRE_LUA = ` local n = redis.call('INCR', KEYS[1]) if n == 1 then redis.call('PEXPIRE', KEYS[1], ARGV[1]) end return n `; const RATE_LIMIT_DEC_LUA = ` local hits = redis.call('DECR', KEYS[1]) if hits < 0 then redis.call('SET', KEYS[1], '0', 'PX', ARGV[1]) hits = 0 end return hits `; const SCAN_BATCH = 200; function parseInfoSections(chunks) { const out = new Map(); for (const chunk of chunks) { for (const line of String(chunk || '').split('\n')) { const trimmed = line.trim(); if (!trimmed || trimmed.startsWith('#')) continue; const sep = trimmed.indexOf(':'); if (sep < 1) continue; out.set(trimmed.slice(0, sep).trim(), trimmed.slice(sep + 1).trim()); } } return out; } function createMemoryBackend() { const store = new Map(); const listeners = new Map(); function live(key) { const e = store.get(key); if (!e) return null; if (e.exp && e.exp <= Date.now()) { store.delete(key); return null; } return e; } return { name: 'memory', async get(key) { const e = live(key); return e ? e.value : undefined; }, async set(key, value, ttlMs) { store.set(key, { value, exp: ttlMs ? Date.now() + ttlMs : 0 }); }, async del(key) { store.delete(key); }, async dropPrefix(prefix) { for (const key of store.keys()) if (key.startsWith(prefix)) store.delete(key); }, async clear() { store.clear(); }, async incr(key, ttlMs) { const e = live(key); if (!e) { store.set(key, { value: 1, exp: ttlMs ? Date.now() + ttlMs : 0 }); return 1; } e.value = (typeof e.value === 'number' ? e.value : 0) + 1; return e.value; }, async rateLimitDec(key, ttlMs) { const e = live(key); if (!e) { store.set(key, { value: 0, exp: Date.now() + ttlMs }); return 0; } e.value = Math.max(0, (typeof e.value === 'number' ? e.value : 0) - 1); return e.value; }, async publish(channel, payload) { const subs = listeners.get(channel); if (!subs || !subs.size) return 0; for (const fn of [...subs]) { try { fn(payload); } catch (e) { console.error('Cache listener failed:', e.message); } } return subs.size; }, on(channel, fn) { if (!listeners.has(channel)) listeners.set(channel, new Set()); const set = listeners.get(channel); set.add(fn); return () => { set.delete(fn); }; }, sweep() { const now = Date.now(); for (const [key, e] of store) if (e.exp && e.exp <= now) store.delete(key); }, }; } function createRedis({ url, prefix } = {}) { const namespace = `${prefix && String(prefix).trim() ? String(prefix).trim() : 'whatido'}:`; const memory = createMemoryBackend(); const handlers = new Map(); const subscribed = new Set(); const stats = { hits: 0, misses: 0, writes: 0, drops: 0, errors: 0, publishes: 0, fallbackOps: 0 }; const CONNECT_TIMEOUT_MS = Math.max(500, parseInt(process.env.REDIS_CONNECT_TIMEOUT_MS || '5000', 10) || 5000); let client = null; let sub = null; let connecting = null; let warned = false; let sweeper = null; let closed = false; function enabled() { return typeof url === 'string' && url.trim() !== ''; } function ready() { return Boolean(client && client.isReady); } function noteError(where, err) { if (closed) return; stats.errors += 1; if (warned) return; warned = true; console.error(`Redis ${where}: ${err && err.message ? err.message : err} — используется in-memory`); } function withTimeout(promise, ms, label) { let timer = null; const guard = new Promise((_, reject) => { timer = setTimeout(() => reject(new Error(`${label} timeout ${ms}ms`)), ms); if (timer.unref) timer.unref(); }); promise.catch(() => {}); return Promise.race([promise, guard]).finally(() => { if (timer) clearTimeout(timer); }); } const fullKey = key => namespace + key; async function via(fn, fallback) { if (!ready()) { stats.fallbackOps += 1; return fallback(); } try { return await fn(client); } catch (e) { noteError('command', e); stats.fallbackOps += 1; return fallback(); } } async function scanDelete(client, target, isGlob) { let cursor = '0'; do { const res = await client.scan(cursor, { MATCH: (isGlob ? target : target + '*'), COUNT: SCAN_BATCH }); cursor = String(res.cursor); if (res.keys.length) await client.del(res.keys); } while (cursor !== '0'); } function deliver(channel, message) { const set = handlers.get(channel); if (!set) return; for (const fn of [...set]) { try { fn(message); } catch (e) { console.error('Redis handler failed:', e.message); } } } async function subscribeChannel(channel) { if (!sub || !sub.isOpen || subscribed.has(channel)) return; try { await sub.subscribe(channel, message => deliver(channel, message)); subscribed.add(channel); } catch (e) { noteError('subscribe', e); } } async function resubscribeAll() { subscribed.clear(); for (const channel of handlers.keys()) await subscribeChannel(channel); } const api = { isEnabled: enabled, isReady: ready, namespace, stats: () => ({ ...stats, enabled: enabled(), ready: ready(), namespace }), async info() { const base = { ...stats, enabled: enabled(), ready: ready(), namespace }; const offline = { ...base, driver: 'memory', version: null, server_uptime_s: null, used_memory_bytes: null, used_memory_human: null, keys: null, }; if (!ready()) return offline; try { const [raw, dbsize, server] = await Promise.all([ client.info('memory'), client.dbSize(), client.info('server'), ]); const fields = parseInfoSections([raw, server]); return { ...base, driver: 'redis', version: fields.get('redis_version') || null, server_uptime_s: fields.has('uptime_in_seconds') ? Number(fields.get('uptime_in_seconds')) || null : null, used_memory_bytes: fields.has('used_memory') ? Number(fields.get('used_memory')) || null : null, used_memory_human: fields.get('used_memory_human') || null, keys: Number(dbsize) || 0, }; } catch (e) { noteError('info', e); return { ...base, driver: 'redis', version: null, server_uptime_s: null, used_memory_bytes: null, used_memory_human: null, keys: null, }; } }, async get(key) { const raw = await via( c => c.get(fullKey(key)), () => memory.get(fullKey(key)) ); if (raw === null || raw === undefined) { stats.misses += 1; return undefined; } stats.hits += 1; try { return JSON.parse(raw); } catch (e) { return raw; } }, async set(key, value, ttlMs) { const payload = JSON.stringify(value === undefined ? null : value); stats.writes += 1; return via( c => (ttlMs ? c.set(fullKey(key), payload, { PX: ttlMs }) : c.set(fullKey(key), payload)), () => memory.set(fullKey(key), value, ttlMs) ); }, async del(key) { return via( c => c.del(fullKey(key)), () => memory.del(fullKey(key)) ); }, async dropPrefix(keyPrefix) { const target = fullKey(keyPrefix); stats.drops += 1; return via( c => scanDelete(c, target), () => memory.dropPrefix(target) ); }, async dropMatch(pattern) { const target = fullKey(pattern); stats.drops += 1; return via( c => scanDelete(c, target, true), () => memory.dropPrefix(target.split('*')[0]) ); }, async clear() { stats.drops += 1; return via( c => scanDelete(c, namespace), () => memory.clear() ); }, async wrap(key, ttlMs, fn) { const hit = await api.get(key); if (hit !== undefined) return hit; const value = await fn(); await api.set(key, value, ttlMs); return value; }, async incr(key, ttlMs) { return via( c => c.eval(INCR_EXPIRE_LUA, { keys: [fullKey(key)], arguments: [String(ttlMs || 0)] }), () => memory.incr(fullKey(key), ttlMs) ); }, rateLimitStore(id, defaultWindowMs) { let windowMs = defaultWindowMs || 60 * 1000; const base = fullKey(`rl:${id}:`); const bucketFor = key => { const bucket = Math.floor(Date.now() / windowMs); return { key: `${base}${key}:${bucket}`, resetTime: (bucket + 1) * windowMs }; }; return { async init(options) { if (options && Number.isFinite(options.windowMs) && options.windowMs > 0) { windowMs = options.windowMs; } }, async increment(key) { const { key: redisKey, resetTime } = bucketFor(key); const ttlMs = windowMs + Math.ceil(windowMs / 10); const hits = await via( c => c.eval(INCR_EXPIRE_LUA, { keys: [redisKey], arguments: [String(ttlMs)] }), () => memory.incr(redisKey, ttlMs) ); return { totalHits: Number(hits), resetTime: new Date(resetTime) }; }, async decrement(key) { const { key: redisKey, resetTime } = bucketFor(key); const ttlMs = Math.max(1, resetTime - Date.now()); await via( c => c.eval(RATE_LIMIT_DEC_LUA, { keys: [redisKey], arguments: [String(ttlMs)] }), () => memory.rateLimitDec(redisKey, ttlMs) ); }, async resetKey(key) { const targets = []; for (let i = 0; i < 2; i += 1) { targets.push(`${base}${key}:${Math.floor(Date.now() / windowMs) - i}`); } return via( c => c.del(targets), async () => { for (const k of targets) await memory.del(k); } ); }, async resetAll() { return via( c => scanDelete(c, base), () => memory.dropPrefix(base) ); }, }; }, async publish(channel, payload) { stats.publishes += 1; const message = typeof payload === 'string' ? payload : JSON.stringify(payload); if (!ready()) { stats.fallbackOps += 1; return memory.publish(channel, message); } try { await client.publish(channel, message); return 1; } catch (e) { noteError('publish', e); stats.fallbackOps += 1; return memory.publish(channel, message); } }, on(channel, fn) { if (!handlers.has(channel)) handlers.set(channel, new Set()); handlers.get(channel).add(fn); memory.on(channel, fn); if (sub && sub.isOpen) subscribeChannel(channel); return () => { const set = handlers.get(channel); if (set) set.delete(fn); }; }, async connect() { if (!enabled()) { console.log('Redis: REDIS_URL не задан, используется in-memory кэш'); return false; } if (connecting) return connecting; if (ready()) return true; closed = false; connecting = (async () => { if (!sweeper) { sweeper = setInterval(() => memory.sweep(), 60 * 1000); sweeper.unref(); } try { client = createClient({ url, socket: { reconnectStrategy: retries => Math.min(200 + retries * 200, 5000) }, }); client.on('error', e => noteError('connection', e)); client.on('ready', () => { warned = false; ensureSubscriber().catch(e => noteError('subscriber', e)); }); await withTimeout(client.connect(), CONNECT_TIMEOUT_MS, 'connect'); await ensureSubscriber(); const pong = await withTimeout(client.ping(), CONNECT_TIMEOUT_MS, 'ping'); console.log(`Redis connected (${pong}), namespace=${namespace}`); return true; } catch (e) { noteError('connect', e); return false; } finally { connecting = null; } })(); return connecting; }, async close() { closed = true; if (sweeper) { clearInterval(sweeper); sweeper = null; } const pending = [sub, client].filter(Boolean); sub = null; client = null; subscribed.clear(); await Promise.all(pending.map(async c => { try { await withTimeout(c.quit(), 2000, 'quit'); } catch (e) { try { c.destroy(); } catch (err) {} } })); }, }; async function ensureSubscriber() { if (closed || !ready()) return; if (!sub) { sub = client.duplicate(); sub.on('error', e => noteError('subscriber', e)); sub.on('ready', () => { resubscribeAll().catch(e => noteError('subscribe', e)); }); } if (!sub.isOpen) await withTimeout(sub.connect(), CONNECT_TIMEOUT_MS, 'subscriber connect'); if (sub.isReady) await resubscribeAll(); } return api; } module.exports = { createRedis, createMemoryBackend };