Files
WhatIDo/redis.js
T
dev 4d3df1de37 fix(redis): разбирать INFO по разделителю, а не по фиксированному смещению
used_memory и used_memory_human вырезались на символ длиннее нужного,
поэтому первая цифра терялась: 1024 превращалось в 24, а "1.34M" — в
".34M". Значение памяти Redis в /api/system-info и в self-тестах было
занижено на порядок.

- parseInfoSections(): разбор INFO по indexOf(':'), пропуск заголовков
  # и пустых строк, работа с любой секцией
- info(): добавлены redis_version и uptime_in_seconds (секция server),
  в том числе в ветках memory и ошибки — с null
2026-09-27 12:46:41 +03:00

465 lines
13 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 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 };