diff --git a/cache-handlers/redis-handler.cjs b/cache-handlers/redis-handler.cjs new file mode 100644 index 0000000..57e88d7 --- /dev/null +++ b/cache-handlers/redis-handler.cjs @@ -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. + } + }, +}; diff --git a/lib/cache/public.ts b/lib/cache/public.ts new file mode 100644 index 0000000..3e3e930 --- /dev/null +++ b/lib/cache/public.ts @@ -0,0 +1,42 @@ +import "server-only"; + +import { cacheLife, cacheTag } from "next/cache"; + +import { getPublicPage, listPublicCharacters } from "@/lib/content/queries"; +import { + characterCacheTag, + dataSourceCacheTag, + DIRECTORY_CACHE_TAG, + pageCacheTag, + vectorCacheTag, +} from "@/lib/cache/tags"; + +export async function getCachedPublicDirectory() { + "use cache: remote"; + cacheLife({ stale: 30, revalidate: 60, expire: 60 * 60 }); + cacheTag(DIRECTORY_CACHE_TAG); + return listPublicCharacters(); +} + +export async function getCachedPublicPage( + characterSlug: string, + pageSlug?: string, +) { + "use cache: remote"; + cacheLife({ stale: 15, revalidate: 60, expire: 60 * 60 }); + const result = await getPublicPage(characterSlug, pageSlug); + cacheTag(DIRECTORY_CACHE_TAG); + if (result.type === "page") { + cacheTag( + pageCacheTag(result.snapshot.page.id), + characterCacheTag(result.snapshot.character.id), + vectorCacheTag( + result.snapshot.page.id, + result.snapshot.page.version, + result.snapshot.versionVector, + ), + ...Object.keys(result.snapshot.versionVector).map(dataSourceCacheTag), + ); + } + return result; +} diff --git a/lib/cache/tags.test.ts b/lib/cache/tags.test.ts new file mode 100644 index 0000000..061dc2a --- /dev/null +++ b/lib/cache/tags.test.ts @@ -0,0 +1,22 @@ +import { describe, expect, it } from "vitest"; + +import { + cacheTagsForOutboxEvent, + DIRECTORY_CACHE_TAG, + vectorCacheTag, +} from "./tags"; + +describe("public cache tags", () => { + it("uses a stable ordered page/data-source version vector", () => { + expect(vectorCacheTag("page", 4, { beta: 2, alpha: 8 })).toBe( + "buzz:vector:page:4:alpha@8,beta@2", + ); + }); + + it("maps outbox events to dependency-aware tags", () => { + expect(cacheTagsForOutboxEvent({ eventType: "directory.updated", aggregateId: "id" })).toEqual([DIRECTORY_CACHE_TAG]); + expect(cacheTagsForOutboxEvent({ eventType: "page.updated", aggregateId: "page" })).toEqual(["buzz:page:page"]); + expect(cacheTagsForOutboxEvent({ eventType: "data-source.updated", aggregateId: "source" })).toEqual(["buzz:data-source:source"]); + expect(cacheTagsForOutboxEvent({ eventType: "admin.updated", aggregateId: "id" })).toEqual([]); + }); +}); diff --git a/lib/cache/tags.ts b/lib/cache/tags.ts new file mode 100644 index 0000000..54224b1 --- /dev/null +++ b/lib/cache/tags.ts @@ -0,0 +1,39 @@ +export const DIRECTORY_CACHE_TAG = "buzz:directory"; + +export function pageCacheTag(pageId: string): string { + return `buzz:page:${pageId}`; +} + +export function characterCacheTag(characterId: string): string { + return `buzz:character:${characterId}`; +} + +export function dataSourceCacheTag(dataSourceId: string): string { + return `buzz:data-source:${dataSourceId}`; +} + +export function vectorCacheTag( + pageId: string, + pageVersion: number, + vector: Record, +): string { + const serialized = Object.entries(vector) + .sort(([left], [right]) => left.localeCompare(right)) + .map(([id, version]) => `${id}@${version}`) + .join(","); + return `buzz:vector:${pageId}:${pageVersion}:${serialized}`; +} + +export type OutboxCacheEvent = { + eventType: string; + aggregateId: string; +}; + +export function cacheTagsForOutboxEvent(event: OutboxCacheEvent): string[] { + if (event.eventType === "directory.updated") return [DIRECTORY_CACHE_TAG]; + if (event.eventType === "page.updated") return [pageCacheTag(event.aggregateId)]; + if (event.eventType === "data-source.updated") { + return [dataSourceCacheTag(event.aggregateId)]; + } + return []; +} diff --git a/lib/content/loaders.ts b/lib/content/loaders.ts index 63bea3b..1206c9f 100644 --- a/lib/content/loaders.ts +++ b/lib/content/loaders.ts @@ -7,24 +7,23 @@ import { getDemoPage, } from "@/lib/content/demo"; import { - getAdminNavigation, - getAdminPage, - getPublicPage, - listPublicCharacters, -} from "@/lib/content/queries"; + getCachedPublicDirectory, + getCachedPublicPage, +} from "@/lib/cache/public"; +import { getAdminNavigation, getAdminPage } from "@/lib/content/queries"; function demoMode(): boolean { return process.env.BUZZ_DEMO_MODE === "true"; } export async function loadPublicDirectory() { - return demoMode() ? demoDirectory : listPublicCharacters(); + return demoMode() ? demoDirectory : getCachedPublicDirectory(); } export async function loadPublicPage(characterSlug: string, pageSlug?: string) { return demoMode() ? getDemoPage(characterSlug, pageSlug) - : getPublicPage(characterSlug, pageSlug); + : getCachedPublicPage(characterSlug, pageSlug); } export async function loadAdminPage(characterSlug: string, pageSlug: string) { diff --git a/next.config.ts b/next.config.ts index e9ffa30..c4a6994 100644 --- a/next.config.ts +++ b/next.config.ts @@ -1,7 +1,21 @@ import type { NextConfig } from "next"; +import { resolve } from "node:path"; const nextConfig: NextConfig = { - /* config options here */ + cacheComponents: true, + cacheHandlers: { + remote: resolve(process.cwd(), "cache-handlers/redis-handler.cjs"), + }, + deploymentId: process.env.NEXT_DEPLOYMENT_ID, + output: "standalone", + experimental: { + serverActions: { + bodySizeLimit: "5mb", + }, + }, + turbopack: { + root: process.cwd(), + }, }; export default nextConfig;