feat(redis): кэш, rate limit, баны IP и pub/sub через Redis

Добавлен сервис 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).
This commit is contained in:
dev
2026-09-26 15:26:00 +03:00
parent e4d58d6525
commit 0e38a280d7
11 changed files with 1098 additions and 74 deletions
+127 -64
View File
@@ -9,6 +9,7 @@ const { createEntryAutoChecker, createPhotoEnhanceWorker } = require('./worker')
const { createZipWriter, renderStudentReport } = require('./student-report');
const { createStorage } = require('./storage');
const { createRedis } = require('./redis');
const https = require('https');
const path = require('path');
@@ -55,39 +56,28 @@ lister.on('notification', (msg) => {
}
});
const cacheStore = new Map();
const cache = createRedis({ url: process.env.REDIS_URL, prefix: process.env.REDIS_PREFIX });
const SETTINGS_TTL_MS = 30 * 1000;
const PUBLIC_TTL_MS = 60 * 1000;
const SHARE_TTL_MS = 60 * 1000;
const STATS_TTL_MS = 15 * 1000;
const SYSTEM_TTL_MS = 30 * 1000;
const SESSION_CACHE_TTL_MS = 30 * 1000;
function cacheGet(key) {
const entry = cacheStore.get(key);
if (!entry) return undefined;
if (entry.exp && entry.exp <= Date.now()) {
cacheStore.delete(key);
return undefined;
}
return entry.value;
async function cacheGet(key) {
return cache.get(key);
}
function cacheSet(key, value, ttlMs) {
cacheStore.set(key, { value, exp: ttlMs ? Date.now() + ttlMs : 0 });
async function cacheSet(key, value, ttlMs) {
return cache.set(key, value, ttlMs);
}
function cacheDrop(prefix) {
for (const key of cacheStore.keys()) {
if (key.startsWith(prefix)) cacheStore.delete(key);
}
function cacheDrop(keyPrefix) {
cache.dropPrefix(keyPrefix).catch(err => console.error('Cache drop failed:', err.message));
}
async function cacheWrap(key, ttlMs, fn) {
const hit = cacheGet(key);
if (hit !== undefined) return hit;
const value = await fn();
cacheSet(key, value, ttlMs);
return value;
return cache.wrap(key, ttlMs, fn);
}
function scopeKey(user) {
@@ -101,31 +91,64 @@ function invalidateSettings() { cacheDrop('setting:'); cacheDrop('share:payload:
function invalidateStudents() { cacheDrop('students:'); }
function invalidateGroups() { cacheDrop('groups:'); cacheDrop('students:'); cacheDrop('share:payload:'); }
function invalidateEntries() { cacheDrop('entries:'); cacheDrop('students:'); cacheDrop('share:payload:'); }
function invalidateSessions() { cacheDrop('session:'); }
const EVENTS_CHANNEL = 'whatido:events';
const AI_WAKE_CHANNEL = 'whatido:wake:ai';
const PHOTO_WAKE_CHANNEL = 'whatido:wake:photo';
function createWorkerBus(channel) {
return {
publish() {
cache.publish(channel, { t: Date.now() }).catch(err => console.error('Publish failed:', err.message));
},
subscribe(fn) {
return cache.on(channel, fn);
},
};
}
const sseClients = new Set();
function broadcastEntryChanged() {
const frame = `event: entries_changed\ndata: ${JSON.stringify({ ts: Date.now() })}\n\n`;
function writeFrame(event, data) {
const frame = `event: ${event}\ndata: ${JSON.stringify(data)}\n\n`;
for (const client of sseClients) {
try { client.write(frame); } catch (e) { sseClients.delete(client); }
}
}
function dispatchEvent(payload) {
if (payload && payload.type === 'ai_status') {
writeFrame('ai_status', {
id: payload.id,
ai_status: payload.status,
ai_error: payload.error || null,
description: payload.description === undefined ? null : payload.description,
description_ai: payload.description_ai === undefined ? null : payload.description_ai,
description_original: payload.description_original === undefined ? null : payload.description_original,
ts: Date.now()
});
return;
}
writeFrame('entries_changed', { ts: Date.now() });
}
function broadcastEntryChanged() {
cache.publish(EVENTS_CHANNEL, { type: 'entries_changed' })
.catch(err => console.error('Publish failed:', err.message));
}
function broadcastAiStatus(payload) {
const data = {
id: payload.id,
ai_status: payload.status,
ai_error: payload.error || null,
description: payload.description === undefined ? null : payload.description,
description_ai: payload.description_ai === undefined ? null : payload.description_ai,
description_original: payload.description_original === undefined ? null : payload.description_original,
ts: Date.now()
};
const frame = `event: ai_status\ndata: ${JSON.stringify(data)}\n\n`;
for (const client of sseClients) {
try { client.write(frame); } catch (e) { sseClients.delete(client); }
}
cache.publish(EVENTS_CHANNEL, { type: 'ai_status', ...payload })
.catch(err => console.error('Publish failed:', err.message));
}
cache.on(EVENTS_CHANNEL, message => {
let payload = null;
try { payload = JSON.parse(message); } catch (e) { return; }
if (payload && payload.type) dispatchEvent(payload);
});
app.get('/api/events', async (req, res) => {
try {
const token = req.headers['x-auth-token'] || req.query.token;
@@ -149,12 +172,12 @@ app.get('/api/events', async (req, res) => {
});
function invalidateShare() { cacheDrop('share:payload:'); }
function invalidateStats() { cacheDrop('stats:'); cacheDrop('dashboard:'); cacheDrop('system-info'); }
function invalidateAll() { cacheStore.clear(); }
function invalidateAll() { cache.clear().catch(err => console.error('Cache clear failed:', err.message)); }
const BAN_TTL_MS = 24 * 60 * 60 * 1000;
const FAIL_WINDOW_MS = 15 * 60 * 1000;
const banMemory = new Map();
const failMemory = new Map();
const banKey = ip => 'ban:' + ip;
const failKey = (kind, ip) => 'fail:' + kind + ':' + ip;
function ipOf(req) {
return String(req.ip || req.socket?.remoteAddress || 'unknown').slice(0, 64);
@@ -166,7 +189,7 @@ async function banIP(req, reason, ms) {
async function banIpAddr(ip, reason, ms, actorReq) {
const until = new Date(Date.now() + ms);
banMemory.set(ip, { reason, banned_until: until.toISOString() });
await cache.set(banKey(ip), { reason, banned_until: until.toISOString() }, ms);
await pool.query(
'INSERT INTO banned_ips (ip, reason, banned_until) VALUES ($1, $2, $3) ON CONFLICT (ip) DO UPDATE SET reason = $2, banned_until = $3',
[ip, reason, until.toISOString()]
@@ -175,40 +198,53 @@ async function banIpAddr(ip, reason, ms, actorReq) {
console.log(`IP banned: ${ip} (${reason})`);
}
function ipGuard(req, res, next) {
const entry = banMemory.get(ipOf(req));
if (entry && new Date(entry.banned_until) > new Date()) {
return res.status(403).json({ error: 'Доступ заблокирован' });
async function unbanIpAddr(ip) {
await cache.del(banKey(ip));
await cache.dropMatch('fail:*:' + ip);
}
async function ipGuard(req, res, next) {
try {
const entry = await cache.get(banKey(ipOf(req)));
if (entry && new Date(entry.banned_until) > new Date()) {
return res.status(403).json({ error: 'Доступ заблокирован' });
}
} catch (e) {
console.error('IP guard failed:', e.message);
}
next();
}
function recordFailure(req, kind, limit, ms) {
const ip = ipOf(req);
const now = Date.now();
let entry = failMemory.get(kind + ':' + ip);
if (!entry || entry.resetAt <= now) {
entry = { count: 0, resetAt: now + FAIL_WINDOW_MS };
failMemory.set(kind + ':' + ip, entry);
}
entry.count += 1;
if (entry.count >= limit) {
failMemory.delete(kind + ':' + ip);
return banIP(req, kind, ms).catch(err => console.error('Ban error:', err));
}
return Promise.resolve();
return cache.incr(failKey(kind, ip), FAIL_WINDOW_MS)
.then(count => {
if (count >= limit) {
return cache.del(failKey(kind, ip))
.then(() => banIP(req, kind, ms))
.catch(err => console.error('Ban error:', err));
}
return null;
})
.catch(err => {
console.error('recordFailure failed:', err.message);
});
}
let seededBans = new Set();
async function loadBans() {
const { rows } = await pool.query('SELECT ip, reason, banned_until FROM banned_ips WHERE banned_until > now()');
const active = new Set();
for (const r of rows) {
active.add(r.ip);
banMemory.set(r.ip, { reason: r.reason, banned_until: r.banned_until });
const ttl = new Date(r.banned_until).getTime() - Date.now();
if (ttl > 0) await cache.set(banKey(r.ip), { reason: r.reason, banned_until: r.banned_until }, ttl);
}
for (const key of banMemory.keys()) {
if (!active.has(key)) banMemory.delete(key);
for (const ip of seededBans) {
if (!active.has(ip)) await cache.del(banKey(ip));
}
seededBans = active;
}
app.set('trust proxy', 'loopback');
@@ -218,6 +254,7 @@ const apiLimiter = rateLimit({
max: 300,
standardHeaders: true,
legacyHeaders: false,
store: cache.rateLimitStore('api', 15 * 60 * 1000),
message: { error: 'Слишком много запросов. Попробуйте позже.' },
});
@@ -226,6 +263,7 @@ const entryLimiter = rateLimit({
max: 10,
standardHeaders: true,
legacyHeaders: false,
store: cache.rateLimitStore('entry', 15 * 60 * 1000),
message: { error: 'Слишком много запросов. Подождите немного.' },
});
@@ -234,6 +272,7 @@ const fileLimiter = rateLimit({
max: 300,
standardHeaders: true,
legacyHeaders: false,
store: cache.rateLimitStore('file', 15 * 60 * 1000),
message: { error: 'Слишком много запросов. Попробуйте позже.' },
});
@@ -421,6 +460,9 @@ function safeUser(u) {
async function loadUserByToken(token) {
if (!token || typeof token !== 'string') return null;
const key = 'session:' + token;
const cached = await cache.get(key);
if (cached !== undefined) return cached;
const { rows } = await pool.query(
`SELECT u.id, u.username, u.name, u.role, u.is_active,
COALESCE(array_agg(ub.branch_id) FILTER (WHERE ub.branch_id IS NOT NULL), '{}') AS branch_ids
@@ -432,6 +474,7 @@ async function loadUserByToken(token) {
[token]
);
if (!rows.length) return null;
await cache.set(key, rows[0], SESSION_CACHE_TTL_MS);
return rows[0];
}
@@ -585,11 +628,11 @@ const adminUpload = multer({
});
async function getSetting(key, def) {
const cached = cacheGet('setting:' + key);
const cached = await cacheGet('setting:' + key);
if (cached !== undefined) return cached;
const { rows } = await pool.query('SELECT value FROM settings WHERE key = $1', [key]);
const value = rows.length ? rows[0].value : def;
cacheSet('setting:' + key, value, SETTINGS_TTL_MS);
await cacheSet('setting:' + key, value, SETTINGS_TTL_MS);
return value;
}
@@ -924,6 +967,7 @@ app.post('/api/auth/login', apiLimiter, async (req, res) => {
app.post('/api/auth/logout', requireAuth, async (req, res) => {
await pool.query('DELETE FROM sessions WHERE token = $1', [req.authToken]);
await cache.del('session:' + req.authToken);
res.json({ ok: true });
});
@@ -947,8 +991,9 @@ app.post('/api/bans', requireAuth, requireAdmin, async (req, res) => {
? req.body.reason.trim().slice(0, 100)
: 'manual';
const hours = Math.min(Math.max(parseInt(req.body?.hours, 10) || 24, 1), 24 * 30);
const bannedUntil = new Date(Date.now() + hours * 60 * 60 * 1000).toISOString();
await banIpAddr(ip, reason, hours * 60 * 60 * 1000, req);
res.json({ ok: true, ip, reason, banned_until: banMemory.get(ip).banned_until });
res.json({ ok: true, ip, reason, banned_until: bannedUntil });
});
app.delete('/api/bans/:ip', requireAuth, requireAdmin, async (req, res) => {
@@ -957,8 +1002,7 @@ app.delete('/api/bans/:ip', requireAuth, requireAdmin, async (req, res) => {
return res.status(400).json({ error: 'Некорректный IP' });
}
await pool.query('DELETE FROM banned_ips WHERE ip = $1', [ip]);
banMemory.delete(ip);
failMemory.forEach((_, key) => { if (key.endsWith(':' + ip)) failMemory.delete(key); });
await unbanIpAddr(ip);
await logAudit(req, 'ip.unban', { ip });
res.json({ ok: true });
});
@@ -1076,6 +1120,7 @@ app.put('/api/users/:id', requireAuth, requireAdmin, async (req, res) => {
await client.query('DELETE FROM sessions WHERE user_id = $1', [id]);
}
await client.query('COMMIT');
invalidateSessions();
await logAudit(req, 'user.update', { id, role: newRole, is_active: isActive });
const fresh = await pool.query(
`SELECT u.id, u.username, u.name, u.role, u.is_active, u.created_at,
@@ -1098,6 +1143,7 @@ app.delete('/api/users/:id', requireAuth, requireAdmin, async (req, res) => {
return res.status(400).json({ error: 'Нельзя удалить самого себя' });
}
await pool.query('DELETE FROM users WHERE id = $1', [req.params.id]);
invalidateSessions();
await logAudit(req, 'user.delete', { id: req.params.id });
res.json({ ok: true });
});
@@ -3907,6 +3953,7 @@ app.get('/api/system-info', requireAdmin, async (_, res) => {
const uploadsCount = usage.count;
const diskInfo = getDiskInfo();
const cacheStats = await cache.info();
return {
database: {
@@ -3942,6 +3989,7 @@ app.get('/api/system-info', requireAdmin, async (_, res) => {
size_bytes: usage.size_bytes,
},
disk: diskInfo,
cache: cacheStats,
};
});
res.json(payload);
@@ -5245,6 +5293,19 @@ const HTTPS_PORT = process.env.HTTPS_PORT || 3443;
process.on('unhandledRejection', (err) => { console.error('Unhandled rejection:', err); });
process.on('uncaughtException', (err) => { console.error('Uncaught exception:', err); });
let shuttingDown = false;
for (const signal of ['SIGTERM', 'SIGINT']) {
process.on(signal, () => {
if (shuttingDown) return;
shuttingDown = true;
console.log(`${signal}: shutting down`);
cache.close()
.catch(() => {})
.finally(() => process.exit(0));
setTimeout(() => process.exit(0), 5000).unref();
});
}
const certPath = path.join(__dirname, 'certs', 'cert.pem');
const keyPath = path.join(__dirname, 'certs', 'key.pem');
@@ -5257,6 +5318,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) {
}
(async () => {
try { await cache.connect(); } catch (err) { console.error('Redis connect:', err); }
try { await ensureBranchesTable(); } catch (err) { console.error('Branches table:', err); }
try { await ensureUsersAndFirstAdmin(); } catch (err) { console.error('Users table:', err); }
try { await ensureAuditTable(); } catch (err) { console.error('Audit table:', err); }
@@ -5285,7 +5347,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) {
setInterval(() => {
try { storage.pruneCache(); } catch (err) { console.error('Cache prune:', err); }
}, 60 * 60 * 1000).unref();
entryAutoChecker = createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl: AI_URL, defaultPrompt: AI_DEFAULT_PROMPT });
entryAutoChecker = createEntryAutoChecker({ pool, getSetting, logAudit, aiUrl: AI_URL, defaultPrompt: AI_DEFAULT_PROMPT, bus: createWorkerBus(AI_WAKE_CHANNEL) });
entryAutoChecker.start();
console.log('AI auto-check worker started');
photoWorker = createPhotoEnhanceWorker({
@@ -5297,6 +5359,7 @@ if (fs.existsSync(certPath) && fs.existsSync(keyPath)) {
photoAiUrl: PHOTO_AI_URL,
uploadsDir: UPLOADS_DIR,
storage,
bus: createWorkerBus(PHOTO_WAKE_CHANNEL),
});
photoWorker.start();
console.log('Photo enhance worker started');