This commit is contained in:
2026-07-20 20:55:09 +08:00 Unverified
parent 52497957f3
commit c31602593a
13 changed files with 453 additions and 31 deletions
+136
View File
@@ -0,0 +1,136 @@
import { createClient } from 'redis';
function positiveInteger(value, fallback, maximum = Number.MAX_SAFE_INTEGER) {
const parsed = Number(value);
return Number.isInteger(parsed) && parsed > 0 ? Math.min(parsed, maximum) : fallback;
}
function disabledCache(status = 'disabled') {
return {
enabled: false,
status,
async remember(_namespace, _key, loader) {
return loader();
},
async invalidate() { return false; },
async close() {}
};
}
export async function createRedisCache({ env = process.env, logger = console, clientFactory = createClient } = {}) {
const url = String(env.REDIS_URL || '').trim();
if (!url) return disabledCache();
const prefix = String(env.REDIS_CACHE_PREFIX || 'exam-information')
.trim()
.replace(/[^a-zA-Z0-9:_-]/g, '-') || 'exam-information';
const defaultTtlSeconds = positiveInteger(env.REDIS_CACHE_TTL_SECONDS, 60, 86400);
const connectTimeout = positiveInteger(env.REDIS_CONNECT_TIMEOUT_MS, 1500, 30000);
const pending = new Map();
let warningReported = false;
const warn = error => {
if (warningReported) return;
warningReported = true;
logger.warn(`Redis 缓存暂不可用,已回源数据库:${error?.message || error}`);
};
const client = clientFactory({
url,
socket: {
connectTimeout,
reconnectStrategy(retries) {
return retries >= 3 ? false : Math.min(100 * 2 ** retries, 1000);
}
}
});
client.on('error', warn);
client.on('ready', () => {
warningReported = false;
});
try {
await client.connect();
} catch (error) {
warn(error);
if (client.isOpen) client.destroy();
return disabledCache('unavailable');
}
const versionKey = namespace => `${prefix}:namespace:${namespace}`;
async function namespaceVersion(namespace) {
const key = versionKey(namespace);
const current = await client.get(key);
if (current) return current;
await client.set(key, '1', { NX: true });
return (await client.get(key)) || '1';
}
return {
get enabled() {
return client.isReady;
},
get status() {
return client.isReady ? 'ready' : 'unavailable';
},
async remember(namespace, key, loader, { ttlSeconds = defaultTtlSeconds } = {}) {
if (!client.isReady) return loader();
try {
const version = await namespaceVersion(namespace);
const cacheKey = `${prefix}:${namespace}:${version}:${key}`;
const cached = await client.get(cacheKey);
if (cached !== null) return JSON.parse(cached);
if (pending.has(cacheKey)) return pending.get(cacheKey);
const loading = Promise.resolve(loader()).then(async value => {
if (client.isReady) {
try {
await client.set(cacheKey, JSON.stringify(value), {
EX: positiveInteger(ttlSeconds, defaultTtlSeconds, 86400)
});
} catch (error) {
warn(error);
}
}
return value;
}).finally(() => pending.delete(cacheKey));
pending.set(cacheKey, loading);
return loading;
} catch (error) {
warn(error);
return loader();
}
},
async invalidate(namespace) {
if (!client.isReady) return false;
try {
await client.incr(versionKey(namespace));
return true;
} catch (error) {
warn(error);
return false;
}
},
async close() {
if (client.isOpen) await client.quit();
}
};
}
export function withCacheInvalidation(database, cache, namespaces = ['public']) {
const resolveNamespaces = typeof namespaces === 'function' ? namespaces : () => namespaces;
return new Proxy(database, {
get(target, property, receiver) {
const value = Reflect.get(target, property, receiver);
if (typeof value !== 'function') return value;
if (property === 'read' || property === 'close') return value.bind(target);
return async (...args) => {
const result = await value.apply(target, args);
const affected = [...new Set(resolveNamespaces(property, args, result) || [])];
await Promise.all(affected.map(namespace => cache.invalidate(namespace)));
return result;
};
}
});
}