Добавлен сервис redis:7-alpine (AOF, requirepass, maxmemory + allkeys-lru, healthcheck, том redis-data, порт только на 127.0.0.1) и абстракция redis.js по образцу storage.js. Переведено на Redis: - кэш ответов API и настроек (было Map в памяти), инвалидация по префиксу через SCAN + DEL; - rate limit для api/entry/file — общие счётчики вместо MemoryStore; - баны IP и счётчики неудачных входа — с TTL, вместо опроса БД каждую минуту; - кэш сессий (30 с) с invalidateSessions() на каждой мутации users/sessions/ user_branches, иначе деактивированный пользователь сохранил бы доступ; - pub/sub для SSE-событий и мгновенного пробуждения фоновых воркеров вместо ожидания цикла опроса БД. Отказоустойчивость: при недоступном Redis все операции уходят в in-memory backend с той же семантикой, приложение стартует и работает без Redis и возвращается в Redis автоматически. Первое подключение ограничено по времени (REDIS_CONNECT_TIMEOUT_MS, 5 с) — node-redis не отклоняет connect() при недоступном сервере, а повторяет попытки бесконечно. Добавлены тесты: redis.selftest.js (в т.ч. поведение при недоступном сервере) и api.smoketest.js (сквозная проверка API, включая инвалидацию кэша и мгновенную смерть сессии после logout).
433 lines
12 KiB
JavaScript
433 lines
12 KiB
JavaScript
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 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 };
|
|
if (!ready()) return { ...base, driver: 'memory', used_memory_bytes: null, used_memory_human: null, keys: null };
|
|
try {
|
|
const [raw, dbsize] = await Promise.all([client.info('memory'), client.dbSize()]);
|
|
let used = null;
|
|
let human = null;
|
|
for (const line of String(raw || '').split('\n')) {
|
|
if (line.startsWith('used_memory:')) used = parseInt(line.slice(13).trim(), 10) || null;
|
|
else if (line.startsWith('used_memory_human:')) human = line.slice(19).trim();
|
|
}
|
|
return {
|
|
...base,
|
|
driver: 'redis',
|
|
used_memory_bytes: used,
|
|
used_memory_human: human,
|
|
keys: Number(dbsize) || 0,
|
|
};
|
|
} catch (e) {
|
|
noteError('info', e);
|
|
return { ...base, driver: 'redis', 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 };
|