183 lines
5.1 KiB
JavaScript
183 lines
5.1 KiB
JavaScript
/* eslint-disable @typescript-eslint/no-require-imports -- Next loads cache handlers as CommonJS modules. */
|
|
const { createHash } = require("node:crypto");
|
|
const RedisModule = require("ioredis");
|
|
|
|
const Redis = RedisModule.default || RedisModule;
|
|
const pendingSets = new Map();
|
|
let redis;
|
|
|
|
function getRedis() {
|
|
if (!process.env.REDIS_URL) {
|
|
throw new Error("REDIS_URL is required for the remote cache handler.");
|
|
}
|
|
if (!redis) {
|
|
redis = new Redis(process.env.REDIS_URL, {
|
|
lazyConnect: true,
|
|
maxRetriesPerRequest: 1,
|
|
enableOfflineQueue: false,
|
|
});
|
|
redis.on("error", () => undefined);
|
|
}
|
|
if (redis.status === "wait") return redis.connect().then(() => redis);
|
|
return Promise.resolve(redis);
|
|
}
|
|
|
|
function prefix() {
|
|
return process.env.REDIS_CACHE_PREFIX || "buzz:next-cache";
|
|
}
|
|
|
|
function entryKey(cacheKey) {
|
|
return `${prefix()}:entry:${createHash("sha256").update(cacheKey).digest("hex")}`;
|
|
}
|
|
|
|
function tagKey() {
|
|
return `${prefix()}:tag-timestamps`;
|
|
}
|
|
|
|
function streamFromBase64(value) {
|
|
const bytes = Buffer.from(value, "base64");
|
|
return new ReadableStream({
|
|
start(controller) {
|
|
controller.enqueue(bytes);
|
|
controller.close();
|
|
},
|
|
});
|
|
}
|
|
|
|
async function readTagStates(client, tags) {
|
|
if (tags.length === 0) return [];
|
|
const values = await client.hmget(tagKey(), ...tags);
|
|
return values.map((value) => {
|
|
if (!value) return {};
|
|
try {
|
|
return JSON.parse(value);
|
|
} catch {
|
|
return {};
|
|
}
|
|
});
|
|
}
|
|
|
|
module.exports = {
|
|
async get(cacheKey, softTags) {
|
|
try {
|
|
const pending = pendingSets.get(cacheKey);
|
|
if (pending) await pending;
|
|
const client = await getRedis();
|
|
const stored = await client.get(entryKey(cacheKey));
|
|
if (!stored) return undefined;
|
|
const entry = JSON.parse(stored);
|
|
const maxAge = process.env.__NEXT_DEV_SERVER
|
|
? Math.max(entry.expire, 30)
|
|
: entry.revalidate;
|
|
if (Date.now() > entry.timestamp + maxAge * 1000) return undefined;
|
|
|
|
const tags = [...new Set([...(entry.tags || []), ...(softTags || [])])];
|
|
const states = await readTagStates(client, tags);
|
|
if (states.some((state) => Number(state.expired || 0) > entry.timestamp)) {
|
|
return undefined;
|
|
}
|
|
const stale = states.some(
|
|
(state) => Number(state.stale || 0) > entry.timestamp,
|
|
);
|
|
return {
|
|
value: streamFromBase64(entry.value),
|
|
tags: entry.tags,
|
|
stale: entry.stale,
|
|
timestamp: entry.timestamp,
|
|
expire: entry.expire,
|
|
revalidate: stale ? -1 : entry.revalidate,
|
|
};
|
|
} catch {
|
|
return undefined;
|
|
}
|
|
},
|
|
|
|
async set(cacheKey, pendingEntry) {
|
|
let resolvePending = () => undefined;
|
|
const pending = new Promise((resolve) => {
|
|
resolvePending = resolve;
|
|
});
|
|
pendingSets.set(cacheKey, pending);
|
|
try {
|
|
const entry = await pendingEntry;
|
|
if (!process.env.__NEXT_DEV_SERVER && entry.expire === 0) return;
|
|
const reader = entry.value.getReader();
|
|
const chunks = [];
|
|
let bytes = 0;
|
|
const maximum = Number(process.env.REDIS_CACHE_MAX_ENTRY_BYTES || 5 * 1024 * 1024);
|
|
try {
|
|
while (true) {
|
|
const chunk = await reader.read();
|
|
if (chunk.done) break;
|
|
const buffer = Buffer.from(chunk.value);
|
|
bytes += buffer.byteLength;
|
|
if (bytes > maximum) return;
|
|
chunks.push(buffer);
|
|
}
|
|
} finally {
|
|
reader.releaseLock();
|
|
}
|
|
const serialized = JSON.stringify({
|
|
value: Buffer.concat(chunks).toString("base64"),
|
|
tags: entry.tags,
|
|
stale: entry.stale,
|
|
timestamp: entry.timestamp,
|
|
expire: entry.expire,
|
|
revalidate: entry.revalidate,
|
|
});
|
|
const ttl = Number.isFinite(entry.expire)
|
|
? Math.max(1, Math.ceil(entry.expire))
|
|
: 86400;
|
|
const client = await getRedis();
|
|
await client.set(entryKey(cacheKey), serialized, "EX", ttl);
|
|
} catch {
|
|
// Cache writes are best effort. The rendered response remains authoritative.
|
|
} finally {
|
|
resolvePending();
|
|
pendingSets.delete(cacheKey);
|
|
}
|
|
},
|
|
|
|
async refreshTags() {
|
|
try {
|
|
const client = await getRedis();
|
|
await client.ping();
|
|
} catch {
|
|
// A Redis outage becomes cache misses, not a request failure.
|
|
}
|
|
},
|
|
|
|
async getExpiration(tags) {
|
|
try {
|
|
const client = await getRedis();
|
|
const states = await readTagStates(client, tags);
|
|
return Math.max(...states.map((state) => Number(state.expired || 0)), 0);
|
|
} catch {
|
|
return 0;
|
|
}
|
|
},
|
|
|
|
async updateTags(tags, durations) {
|
|
try {
|
|
const client = await getRedis();
|
|
const now = Date.now();
|
|
const pipeline = client.pipeline();
|
|
for (const tag of tags) {
|
|
const state = durations
|
|
? {
|
|
stale: now,
|
|
expired:
|
|
durations.expire === undefined
|
|
? undefined
|
|
: now + durations.expire * 1000,
|
|
}
|
|
: { expired: now };
|
|
pipeline.hset(tagKey(), tag, JSON.stringify(state));
|
|
}
|
|
await pipeline.exec();
|
|
} catch {
|
|
// Outbox-driven invalidations retry independently.
|
|
}
|
|
},
|
|
};
|