feat(cache): share public snapshots through Redis
This commit is contained in:
@@ -0,0 +1,182 @@
|
||||
/* 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.
|
||||
}
|
||||
},
|
||||
};
|
||||
Reference in New Issue
Block a user