182 lines
6.4 KiB
JavaScript
182 lines
6.4 KiB
JavaScript
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', { ttlSeconds = 60, maxEntries = 200 } = {}) {
|
|
const values = new Map();
|
|
const pending = new Map();
|
|
const generations = new Map();
|
|
|
|
const cacheKey = (namespace, key) => `${namespace}:${key}`;
|
|
const generation = namespace => generations.get(namespace) || 0;
|
|
const prune = () => {
|
|
const now = Date.now();
|
|
for (const [key, entry] of values) if (entry.expiresAt <= now) values.delete(key);
|
|
while (values.size > maxEntries) values.delete(values.keys().next().value);
|
|
};
|
|
|
|
return {
|
|
enabled: false,
|
|
status,
|
|
async remember(namespace, key, loader, options = {}) {
|
|
const fullKey = cacheKey(namespace, key);
|
|
const cached = values.get(fullKey);
|
|
if (cached && cached.expiresAt > Date.now()) return cached.value;
|
|
if (cached) values.delete(fullKey);
|
|
const startedGeneration = generation(namespace);
|
|
const active = pending.get(fullKey);
|
|
if (active?.generation === startedGeneration) return active.promise;
|
|
|
|
const loading = Promise.resolve(loader()).then(value => {
|
|
if (generation(namespace) === startedGeneration) {
|
|
const lifetime = positiveInteger(options.ttlSeconds, ttlSeconds, 86400);
|
|
values.set(fullKey, { value, expiresAt: Date.now() + lifetime * 1000 });
|
|
prune();
|
|
}
|
|
return value;
|
|
}).finally(() => {
|
|
if (pending.get(fullKey)?.promise === loading) pending.delete(fullKey);
|
|
});
|
|
pending.set(fullKey, { generation: startedGeneration, promise: loading });
|
|
return loading;
|
|
},
|
|
async invalidate(namespace) {
|
|
generations.set(namespace, generation(namespace) + 1);
|
|
const prefix = `${namespace}:`;
|
|
for (const key of values.keys()) if (key.startsWith(prefix)) values.delete(key);
|
|
// Preserve the public meaning of this return value: no Redis namespace
|
|
// was refreshed, even though the local fallback was invalidated.
|
|
return false;
|
|
},
|
|
async close() {
|
|
values.clear();
|
|
pending.clear();
|
|
}
|
|
};
|
|
}
|
|
|
|
export async function createRedisCache({ env = process.env, logger = console, clientFactory = createClient } = {}) {
|
|
const url = String(env.REDIS_URL || '').trim();
|
|
const defaultTtlSeconds = positiveInteger(env.REDIS_CACHE_TTL_SECONDS, 60, 86400);
|
|
const localMaxEntries = positiveInteger(env.LOCAL_CACHE_MAX_ENTRIES, 200, 5000);
|
|
if (!url) return disabledCache('disabled', { ttlSeconds: defaultTtlSeconds, maxEntries: localMaxEntries });
|
|
const fallback = disabledCache('unavailable', { ttlSeconds: defaultTtlSeconds, maxEntries: localMaxEntries });
|
|
|
|
const prefix = String(env.REDIS_CACHE_PREFIX || 'exam-information')
|
|
.trim()
|
|
.replace(/[^a-zA-Z0-9:_-]/g, '-') || 'exam-information';
|
|
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 fallback;
|
|
}
|
|
|
|
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 fallback.remember(namespace, key, loader, { ttlSeconds });
|
|
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 fallback.remember(namespace, key, loader, { ttlSeconds });
|
|
}
|
|
},
|
|
async invalidate(namespace) {
|
|
await fallback.invalidate(namespace);
|
|
if (!client.isReady) return false;
|
|
try {
|
|
await client.incr(versionKey(namespace));
|
|
return true;
|
|
} catch (error) {
|
|
warn(error);
|
|
return false;
|
|
}
|
|
},
|
|
async close() {
|
|
await fallback.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;
|
|
};
|
|
}
|
|
});
|
|
}
|