diff --git a/.env.example b/.env.example index ebed6b34..8668f6b8 100644 --- a/.env.example +++ b/.env.example @@ -11,6 +11,12 @@ DATABASE_CONNECT_TIMEOUT_MS=5000 # --- REDIS (Lightning Fast Caching & Sessions) --- REDIS_URL=redis://127.0.0.1:6379?connect_timeout=2 REDIS_CACHE_TTL_DEFAULT=7200 +# In-process cache entries kept per instance, evicted least-recently-used. Raise +# it if hot keys are evicted while memory headroom remains (default 2000). +CACHE_MEMORY_MAX_ENTRIES=2000 +# Renders kept per imaging cache directory, counted as .img/.json pairs. A sweep +# every 5 minutes brings an over-budget directory back to 90% of this (default 20000). +IMAGING_CACHE_MAX_ENTRIES=20000 # --- CORE RUNTIME & PERFORMANCE FLAGS --- NODE_ENV=production diff --git a/src/app/(site)/leaderboard/page.tsx b/src/app/(site)/leaderboard/page.tsx index 238b3db5..a6b87d5c 100644 --- a/src/app/(site)/leaderboard/page.tsx +++ b/src/app/(site)/leaderboard/page.tsx @@ -56,16 +56,20 @@ type Row = { username: string; look: string; value: number }; async function loadCreditsRows(): Promise { try { - const users = await cached("lb_credits", 60_000, () => - db - .select({ - username: User.username, - look: User.look, - credits: User.credits, - }) - .from(User) - .orderBy(desc(User.credits)) - .limit(20), + const users = await cached( + "lb_credits", + 60_000, + () => + db + .select({ + username: User.username, + look: User.look, + credits: User.credits, + }) + .from(User) + .orderBy(desc(User.credits)) + .limit(20), + { staleMs: 120000 }, ); return users.map((u) => ({ username: u.username, @@ -79,18 +83,22 @@ async function loadCreditsRows(): Promise { async function loadCurrencyRows(type: number): Promise { try { - return await cached(`lb_currency_${type}`, 60_000, () => - db - .select({ - username: User.username, - look: User.look, - value: UsersCurrency.amount, - }) - .from(UsersCurrency) - .innerJoin(User, eq(UsersCurrency.userId, User.id)) - .where(eq(UsersCurrency.type, type)) - .orderBy(desc(UsersCurrency.amount)) - .limit(20), + return await cached( + `lb_currency_${type}`, + 60_000, + () => + db + .select({ + username: User.username, + look: User.look, + value: UsersCurrency.amount, + }) + .from(UsersCurrency) + .innerJoin(User, eq(UsersCurrency.userId, User.id)) + .where(eq(UsersCurrency.type, type)) + .orderBy(desc(UsersCurrency.amount)) + .limit(20), + { staleMs: 120000 }, ); } catch { return []; @@ -102,17 +110,21 @@ async function loadSettingsRows( ): Promise { try { const column = UsersSettings[field]; - return await cached(`lb_settings_${field}`, 60_000, () => - db - .select({ - username: User.username, - look: User.look, - value: column, - }) - .from(UsersSettings) - .innerJoin(User, eq(UsersSettings.userId, User.id)) - .orderBy(desc(column)) - .limit(20), + return await cached( + `lb_settings_${field}`, + 60_000, + () => + db + .select({ + username: User.username, + look: User.look, + value: column, + }) + .from(UsersSettings) + .innerJoin(User, eq(UsersSettings.userId, User.id)) + .orderBy(desc(column)) + .limit(20), + { staleMs: 120000 }, ); } catch { return []; diff --git a/src/app/(site)/page.tsx b/src/app/(site)/page.tsx index 878d9ce9..31396a0c 100644 --- a/src/app/(site)/page.tsx +++ b/src/app/(site)/page.tsx @@ -181,40 +181,60 @@ async function getHotelData() { .where(eq(User.online, "1")) .then((rows) => rows[0]?.total ?? 0), ).catch(publicReadFailure("home.online")), - cached("total_users", 300_000, () => - db - .select({ total: count() }) - .from(User) - .then((rows) => rows[0]?.total ?? 0), + cached( + "total_users", + 300_000, + () => + db + .select({ total: count() }) + .from(User) + .then((rows) => rows[0]?.total ?? 0), + { staleMs: 300000 }, ).catch(publicReadFailure("home.users")), - cached("total_rooms", 300_000, () => - db - .select({ total: count() }) - .from(Rooms) - .then((rows) => rows[0]?.total ?? 0), + cached( + "total_rooms", + 300_000, + () => + db + .select({ total: count() }) + .from(Rooms) + .then((rows) => rows[0]?.total ?? 0), + { staleMs: 300000 }, ).catch(publicReadFailure("home.rooms")), - cached("total_photos", 300_000, () => - db - .select({ total: count() }) - .from(CameraWeb) - .then((rows) => rows[0]?.total ?? 0), + cached( + "total_photos", + 300_000, + () => + db + .select({ total: count() }) + .from(CameraWeb) + .then((rows) => rows[0]?.total ?? 0), + { staleMs: 300000 }, ).catch(publicReadFailure("home.photos-count")), getNewsList(4, { throwOnError: true }).catch( publicReadFailure("home.news"), ), - cached("home_online_users", 15_000, () => - db - .select({ username: User.username, look: User.look }) - .from(User) - .where(eq(User.online, "1")) - .limit(12), + cached( + "home_online_users", + 15_000, + () => + db + .select({ username: User.username, look: User.look }) + .from(User) + .where(eq(User.online, "1")) + .limit(12), + { staleMs: 30000 }, ).catch(publicReadFailure("home.recent-users")), - cached("home_recent_photos", 60_000, () => - db - .select({ id: CameraWeb.id, url: CameraWeb.url }) - .from(CameraWeb) - .orderBy(desc(CameraWeb.timestamp)) - .limit(4), + cached( + "home_recent_photos", + 60_000, + () => + db + .select({ id: CameraWeb.id, url: CameraWeb.url }) + .from(CameraWeb) + .orderBy(desc(CameraWeb.timestamp)) + .limit(4), + { staleMs: 120000 }, ).catch(publicReadFailure("home.photos")), ]); diff --git a/src/app/(site)/photos/page.tsx b/src/app/(site)/photos/page.tsx index 77dac002..581a2a8d 100644 --- a/src/app/(site)/photos/page.tsx +++ b/src/app/(site)/photos/page.tsx @@ -38,8 +38,16 @@ export default async function PhotosPage() { let photos: Photo[] | null = []; try { - photos = await cached("photos:grid", 60_000, async () => - db.select().from(CameraWeb).orderBy(desc(CameraWeb.timestamp)).limit(48), + photos = await cached( + "photos:grid", + 60_000, + async () => + db + .select() + .from(CameraWeb) + .orderBy(desc(CameraWeb.timestamp)) + .limit(48), + { staleMs: 120000 }, ); } catch (error) { photos = publicReadFailure("photos")(error); diff --git a/src/app/(site)/rankings/page.tsx b/src/app/(site)/rankings/page.tsx index 204cc126..4982befb 100644 --- a/src/app/(site)/rankings/page.tsx +++ b/src/app/(site)/rankings/page.tsx @@ -40,18 +40,22 @@ export default async function RankingsPage() { const t = await getTranslations("pages.rankings"); let users: TopUser[] | null = []; try { - users = await cached("rankings:top", 60_000, async () => - db - .select({ - username: User.username, - look: User.look, - credits: User.credits, - rank: User.rank, - online: User.online, - }) - .from(User) - .orderBy(desc(User.credits)) - .limit(12), + users = await cached( + "rankings:top", + 60_000, + async () => + db + .select({ + username: User.username, + look: User.look, + credits: User.credits, + rank: User.rank, + online: User.online, + }) + .from(User) + .orderBy(desc(User.credits)) + .limit(12), + { staleMs: 120000 }, ); } catch (error) { users = publicReadFailure("rankings")(error); diff --git a/src/app/(site)/register/page.tsx b/src/app/(site)/register/page.tsx index 8cd13634..81a751dc 100644 --- a/src/app/(site)/register/page.tsx +++ b/src/app/(site)/register/page.tsx @@ -31,19 +31,27 @@ export default async function RegisterPage() { .where(eq(User.online, "1")) .then((rows) => rows[0]?.total ?? 0), ).catch(() => 0), - cached("register_online_users", 10_000, () => - db - .select({ username: User.username, look: User.look }) - .from(User) - .where(eq(User.online, "1")) - .limit(8), + cached( + "register_online_users", + 10_000, + () => + db + .select({ username: User.username, look: User.look }) + .from(User) + .where(eq(User.online, "1")) + .limit(8), + { staleMs: 30000 }, ).catch(() => []), - cached("register_latest_users", 30_000, () => - db - .select({ username: User.username, look: User.look }) - .from(User) - .orderBy(desc(User.accountCreated)) - .limit(8), + cached( + "register_latest_users", + 30_000, + () => + db + .select({ username: User.username, look: User.look }) + .from(User) + .orderBy(desc(User.accountCreated)) + .limit(8), + { staleMs: 60000 }, ).catch(() => []), ]); diff --git a/src/app/api/admin/devops/cache/route.ts b/src/app/api/admin/devops/cache/route.ts new file mode 100644 index 00000000..34725582 --- /dev/null +++ b/src/app/api/admin/devops/cache/route.ts @@ -0,0 +1,15 @@ +import { NextResponse } from "next/server"; +import { withAdmin } from "@/lib/api-handler"; +import { cacheBudget } from "@/lib/cache"; +import { readCacheStats } from "@/lib/cache-stats"; +import { PERMS } from "@/lib/permissions"; + +/** + * Cache health report: per-key hit/miss/stale/error/eviction counters plus the + * in-process budget, so a cache that is silently doing nothing is visible + * instead of only showing up as "the site got slower". + */ +export const GET = withAdmin({ permission: PERMS.DEVOPS_VIEW }, async () => { + const stats = await readCacheStats(); + return NextResponse.json({ ...stats, budget: cacheBudget() }); +}); diff --git a/src/app/api/badges/leaderboard/route.ts b/src/app/api/badges/leaderboard/route.ts index d787c71d..2dc33086 100644 --- a/src/app/api/badges/leaderboard/route.ts +++ b/src/app/api/badges/leaderboard/route.ts @@ -362,6 +362,7 @@ export async function GET(req: Request) { return cacheSafe({ badgeStats, totalBadges, achievementLevel, rarity }); }, + { staleMs: 120_000 }, ); // Per-viewer rank entries (only when a user is signed in). diff --git a/src/app/api/guilds/[id]/route.ts b/src/app/api/guilds/[id]/route.ts index 3187c178..dab8f9cd 100644 --- a/src/app/api/guilds/[id]/route.ts +++ b/src/app/api/guilds/[id]/route.ts @@ -56,6 +56,7 @@ export async function GET( const { userId, ...rest } = guild; return cacheSafe({ data: { ...rest, ownerId: userId, memberCount } }); }, + { staleMs: 120_000 }, ); if (data === null) { diff --git a/src/app/api/guilds/route.ts b/src/app/api/guilds/route.ts index b7cce013..d276f6f1 100644 --- a/src/app/api/guilds/route.ts +++ b/src/app/api/guilds/route.ts @@ -50,6 +50,7 @@ export async function GET(req: Request) { }, }); }, + { staleMs: 120_000 }, ); return apiJson(data); diff --git a/src/app/api/leaderboard/route.ts b/src/app/api/leaderboard/route.ts index 09c795c6..e9535714 100644 --- a/src/app/api/leaderboard/route.ts +++ b/src/app/api/leaderboard/route.ts @@ -84,6 +84,9 @@ export async function GET(req: Request) { ? await loadCreditsRows() : // eslint-disable-next-line security/detect-object-injection -- type validated to "diamonds"|"duckets" await loadCurrencyRows(CURRENCY_TYPE[type]), + // Rankings are re-sorted on every read; serving the previous order for a + // minute keeps a busy leaderboard off the database. + { staleMs: 60 }, ); return apiJson({ type, data: rows }, { status: 200 }); diff --git a/src/app/api/online/count/route.ts b/src/app/api/online/count/route.ts index a684a831..8d299816 100644 --- a/src/app/api/online/count/route.ts +++ b/src/app/api/online/count/route.ts @@ -10,13 +10,21 @@ import { db, User } from "@/lib/db"; export async function GET(_req: Request) { try { - const result = await cached("online_count", 10_000, async () => { - const [row] = await db - .select({ total: count() }) - .from(User) - .where(eq(User.online, "1")); - return row?.total ?? 0; - }); + // Polled by every page view and the in-game SSE stream, so a count that + // is one interval behind is worthless as "wrong" but a blocking COUNT(*) + // on each expiry is not. The grace window keeps the origin off the path. + const result = await cached( + "online_count", + 10_000, + async () => { + const [row] = await db + .select({ total: count() }) + .from(User) + .where(eq(User.online, "1")); + return row?.total ?? 0; + }, + { staleMs: 15_000 }, + ); return apiJson({ count: result }); } catch { return apiJson({ count: 0 }); diff --git a/src/app/api/online/count/stream/route.ts b/src/app/api/online/count/stream/route.ts index ea9b590f..eddd9ef4 100644 --- a/src/app/api/online/count/stream/route.ts +++ b/src/app/api/online/count/stream/route.ts @@ -13,13 +13,19 @@ import { rcon } from "@/lib/services/rcon"; async function fetchOnlineCount(): Promise { try { - return await cached("online_count", 10_000, async () => { - const [row] = await db - .select({ total: count() }) - .from(User) - .where(eq(User.online, "1")); - return row?.total ?? 0; - }); + return await cached( + "online_count", + 10_000, + async () => { + const [row] = await db + .select({ total: count() }) + .from(User) + .where(eq(User.online, "1")); + return row?.total ?? 0; + }, + // A subscriber must never wait on a COUNT(*) to get its next event. + { staleMs: 15_000 }, + ); } catch { return 0; } @@ -30,8 +36,13 @@ async function fetchOnlineCount(): Promise { // degradation on failure so the stream never errors out. async function fetchEmulatorStatus(): Promise { try { - return await cached("emulator_rcon_status", 15_000, () => - rcon.send("ping", null).catch(() => false), + return await cached( + "emulator_rcon_status", + 15_000, + () => rcon.send("ping", null).catch(() => false), + // Reachability is reported, never acted on, so a previous answer is + // good enough while the ping is re-run. + { staleMs: 15_000 }, ); } catch { return false; diff --git a/src/app/api/online/route.ts b/src/app/api/online/route.ts index ef3309e9..2d183fb5 100644 --- a/src/app/api/online/route.ts +++ b/src/app/api/online/route.ts @@ -10,12 +10,18 @@ import { db, User } from "@/lib/db"; export async function GET(_req: Request) { try { - const users = await cached("online_users", 10_000, async () => - db - .select({ username: User.username, look: User.look }) - .from(User) - .where(eq(User.online, "1")) - .limit(100), + // A briefly stale online list is harmless; blocking the first request + // after the TTL on a COUNT + SELECT is not. + const users = await cached( + "online_users", + 10_000, + async () => + db + .select({ username: User.username, look: User.look }) + .from(User) + .where(eq(User.online, "1")) + .limit(100), + { staleMs: 15_000 }, ); return apiJson({ users }); diff --git a/src/app/api/photos/route.ts b/src/app/api/photos/route.ts index f1fd212c..4ec59743 100644 --- a/src/app/api/photos/route.ts +++ b/src/app/api/photos/route.ts @@ -44,6 +44,7 @@ export async function GET(req: Request) { }, }); }, + { staleMs: 120_000 }, ); return apiJson(data); diff --git a/src/app/api/radio/current-dj/route.ts b/src/app/api/radio/current-dj/route.ts index f316f99a..1ce83a3b 100644 --- a/src/app/api/radio/current-dj/route.ts +++ b/src/app/api/radio/current-dj/route.ts @@ -38,6 +38,7 @@ export async function GET(_req: Request) { return { username: user.username, look: user.look }; }, + { staleMs: 30_000 }, ); return apiJson({ dj }); diff --git a/src/app/api/radio/points/leaderboard/route.ts b/src/app/api/radio/points/leaderboard/route.ts index 333b7b73..048342d5 100644 --- a/src/app/api/radio/points/leaderboard/route.ts +++ b/src/app/api/radio/points/leaderboard/route.ts @@ -48,6 +48,7 @@ export async function GET(_req: Request) { }; }); }, + { staleMs: 120_000 }, ); return apiJson({ data }); diff --git a/src/app/api/radio/shouts/route.ts b/src/app/api/radio/shouts/route.ts index c81c5b6a..15a66297 100644 --- a/src/app/api/radio/shouts/route.ts +++ b/src/app/api/radio/shouts/route.ts @@ -58,6 +58,7 @@ export async function GET(_req: Request) { }; }); }, + { staleMs: 10_000 }, ); return apiJson({ shouts: data }); diff --git a/src/app/api/shop/categories/route.ts b/src/app/api/shop/categories/route.ts index 1c2ba1be..da368113 100644 --- a/src/app/api/shop/categories/route.ts +++ b/src/app/api/shop/categories/route.ts @@ -26,6 +26,7 @@ export async function GET(_req: Request) { asc(WebsiteShopCategories.name), ), ), + { staleMs: 120_000 }, ); return apiJson({ data }); diff --git a/src/app/api/shop/route.ts b/src/app/api/shop/route.ts index cbf51546..c02fb34b 100644 --- a/src/app/api/shop/route.ts +++ b/src/app/api/shop/route.ts @@ -71,6 +71,7 @@ export async function GET(req: Request) { }, }); }, + { staleMs: 120_000 }, ); return apiJson(data); diff --git a/src/app/api/staff/route.ts b/src/app/api/staff/route.ts index ad4036bc..e4a52d7f 100644 --- a/src/app/api/staff/route.ts +++ b/src/app/api/staff/route.ts @@ -12,18 +12,24 @@ export async function GET(_req: Request) { const minStaffRank = Number(await siteSettings.get("min_staff_rank", "7")) || 7; - const staff = await redisCache(apiCacheKey("staff"), 300, () => - db - .select({ - username: User.username, - look: User.look, - rank: User.rank, - motto: User.motto, - }) - .from(User) - .where(gte(User.rank, minStaffRank)) - .orderBy(desc(User.rank), asc(User.username)) - .limit(100), + const staff = await redisCache( + apiCacheKey("staff"), + 300, + () => + db + .select({ + username: User.username, + look: User.look, + rank: User.rank, + motto: User.motto, + }) + .from(User) + .where(gte(User.rank, minStaffRank)) + .orderBy(desc(User.rank), asc(User.username)) + .limit(100), + // Staff badges are rendered on every public page; a roster that is a + // few minutes behind must never turn into a query on each visitor. + { staleMs: 300 }, ); return apiJson({ data: staff }, { status: 200 }); diff --git a/src/app/api/teams/route.ts b/src/app/api/teams/route.ts index 4eb58ba2..5e5baa4a 100644 --- a/src/app/api/teams/route.ts +++ b/src/app/api/teams/route.ts @@ -8,20 +8,24 @@ import { apiCacheKey, cacheSafe, redisCache } from "@/lib/redis-cache"; export async function GET(_req: Request) { try { - const rows = await redisCache(apiCacheKey("teams"), 300, async () => - cacheSafe( - await db - .select({ - id: WebsiteTeams.id, - rankName: WebsiteTeams.rankName, - badge: WebsiteTeams.badge, - jobDescription: WebsiteTeams.jobDescription, - staffColor: WebsiteTeams.staffColor, - }) - .from(WebsiteTeams) - .where(eq(WebsiteTeams.hiddenRank, false)) - .orderBy(asc(WebsiteTeams.id)), - ), + const rows = await redisCache( + apiCacheKey("teams"), + 300, + async () => + cacheSafe( + await db + .select({ + id: WebsiteTeams.id, + rankName: WebsiteTeams.rankName, + badge: WebsiteTeams.badge, + jobDescription: WebsiteTeams.jobDescription, + staffColor: WebsiteTeams.staffColor, + }) + .from(WebsiteTeams) + .where(eq(WebsiteTeams.hiddenRank, false)) + .orderBy(asc(WebsiteTeams.id)), + ), + { staleMs: 600_000 }, ); // apiJson serialises BigInt ids → string automatically. diff --git a/src/app/api/values/[id]/route.ts b/src/app/api/values/[id]/route.ts index d4ca7884..dd3ab18e 100644 --- a/src/app/api/values/[id]/route.ts +++ b/src/app/api/values/[id]/route.ts @@ -76,6 +76,7 @@ export async function GET( return cacheSafe({ data: { ...value, category } }); }, + { staleMs: 300_000 }, ); if (data === null) { diff --git a/src/app/api/values/categories/route.ts b/src/app/api/values/categories/route.ts index 0b925bb5..eb2e4054 100644 --- a/src/app/api/values/categories/route.ts +++ b/src/app/api/values/categories/route.ts @@ -26,6 +26,7 @@ export async function GET(_req: Request) { asc(WebsiteRareValueCategories.name), ), ), + { staleMs: 600_000 }, ); return apiJson({ data }); diff --git a/src/app/api/values/route.ts b/src/app/api/values/route.ts index 90cc8115..1c8551b3 100644 --- a/src/app/api/values/route.ts +++ b/src/app/api/values/route.ts @@ -61,6 +61,7 @@ export async function GET(req: Request) { }, }); }, + { staleMs: 300_000 }, ); return apiJson(data); diff --git a/src/instrumentation.ts b/src/instrumentation.ts index 5bdf6433..b2506325 100644 --- a/src/instrumentation.ts +++ b/src/instrumentation.ts @@ -39,4 +39,10 @@ export async function register() { void sweepExpiredCloudflareBlocks(); }, 30000); cloudflareSweep.unref(); + + // A deploy starts with an empty in-process cache, so prime the hot public + // keys before traffic arrives instead of letting the first visitors each pay + // for a miss. Never awaited: it must not delay the server becoming ready. + const { warmPublicCaches } = await import("@/lib/services/cache-warmup"); + void warmPublicCaches(); } diff --git a/src/lib/auth/jwt-version-cache.ts b/src/lib/auth/jwt-version-cache.ts index f5ce6a13..7388b1e2 100644 --- a/src/lib/auth/jwt-version-cache.ts +++ b/src/lib/auth/jwt-version-cache.ts @@ -8,6 +8,24 @@ const MEMORY_TTL_MS = 60_000; const REDIS_TTL_SEC = 60; const memory = new Map(); +// One entry per authenticated user, so without a cap a long-lived process grows +// this map by every account that ever signs in. Entries are tiny, so a generous +// bound costs almost nothing and keeps the leak impossible. +const MAX_MEMORY_ENTRIES = 20_000; + +function remember(userId: number, version: number, now: number): void { + if (memory.size >= MAX_MEMORY_ENTRIES) { + for (const [id, entry] of memory) { + if (entry.expiresAt <= now) memory.delete(id); + } + if (memory.size >= MAX_MEMORY_ENTRIES) { + const oldest = memory.keys().next().value; + if (oldest !== undefined) memory.delete(oldest); + } + } + memory.set(userId, { version, expiresAt: now + MEMORY_TTL_MS }); +} + function redisKey(userId: number): string { return `jwtver:${userId}`; } @@ -31,7 +49,7 @@ export async function getCachedJwtVersion( if (raw !== null && raw !== undefined) { const version = Number.parseInt(raw, 10); if (Number.isFinite(version)) { - memory.set(userId, { version, expiresAt: now + MEMORY_TTL_MS }); + remember(userId, version, now); return version; } } @@ -48,7 +66,7 @@ export async function getCachedJwtVersion( .limit(1); if (!row) return null; const version = row.websiteJwtVersion; - memory.set(userId, { version, expiresAt: now + MEMORY_TTL_MS }); + remember(userId, version, now); if (redis) { try { await redis.setex(redisKey(userId), REDIS_TTL_SEC, String(version)); diff --git a/src/lib/cache-invalidation.test.ts b/src/lib/cache-invalidation.test.ts new file mode 100644 index 00000000..d8b72738 --- /dev/null +++ b/src/lib/cache-invalidation.test.ts @@ -0,0 +1,105 @@ +// @ts-nocheck +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; + +const state = vi.hoisted(() => ({ + duplicateCalls: 0, + subscribeCalls: 0, + published: [] as string[], + // Simulates another process invalidating a key: delivers a pub/sub message + // to whichever subscriber this module registered. + deliver: null as ((channel: string, message: string) => void) | null, +})); + +vi.mock("@/lib/redis", () => ({ + redis: { + status: "ready", + publish: async (channel: string, key: string) => { + state.published.push(`${channel}|${key}`); + return 1; + }, + duplicate: () => { + state.duplicateCalls++; + return { + on: (event: string, handler: (...args: unknown[]) => void) => { + if (event === "message") state.deliver = handler; + }, + subscribe: async () => { + state.subscribeCalls++; + return 1; + }, + }; + }, + }, +})); + +const { onCacheInvalidated, publishCacheInvalidation } = await import( + "./cache-invalidation" +); + +// The module registers one subscriber for the whole file, so its message handler +// must stay wired between tests; only the subscriptions a test creates are torn +// down, otherwise handlers leak from one test into the next. +const subscriptions: (() => void)[] = []; +function track(handler: (key: string) => void): void { + subscriptions.push(onCacheInvalidated(handler)); +} + +beforeEach(() => { + state.duplicateCalls = 0; + state.subscribeCalls = 0; + state.published = []; +}); + +afterEach(() => { + while (subscriptions.length) subscriptions.pop()?.(); + vi.unstubAllEnvs(); +}); + +describe("cache invalidation transport", () => { + it("delivers a key invalidated by another process to local handlers", () => { + const seen: string[] = []; + track((key) => seen.push(key)); + state.deliver?.("cache:invalidate", "api:staff"); + expect(seen).toEqual(["api:staff"]); + }); + + it("ignores messages on other channels", () => { + const seen: string[] = []; + track((key) => seen.push(key)); + state.deliver?.("some:other:channel", "api:staff"); + expect(seen).toEqual([]); + }); + + it("stops delivering after unsubscribe", () => { + const seen: string[] = []; + onCacheInvalidated((key) => seen.push(key))(); + state.deliver?.("cache:invalidate", "api:staff"); + expect(seen).toEqual([]); + }); + + it("keeps delivering to the remaining handlers when one throws", () => { + const seen: string[] = []; + track(() => { + throw new Error("handler blew up"); + }); + track((key) => seen.push(key)); + expect(() => state.deliver?.("cache:invalidate", "k")).not.toThrow(); + expect(seen).toEqual(["k"]); + }); + + it("publishes the invalidated key on the shared channel", async () => { + await publishCacheInvalidation("api:staff"); + expect(state.published).toEqual(["cache:invalidate|api:staff"]); + }); + + it("opens no connection during a production build", async () => { + // A live subscriber socket would keep `next build` from ever exiting. + vi.stubEnv("NEXT_PHASE", "phase-production-build"); + vi.resetModules(); + const build = await import("./cache-invalidation"); + build.onCacheInvalidated(() => {}); + await build.publishCacheInvalidation("api:staff"); + expect(state.duplicateCalls).toBe(0); + expect(state.subscribeCalls).toBe(0); + }); +}); diff --git a/src/lib/cache-invalidation.ts b/src/lib/cache-invalidation.ts new file mode 100644 index 00000000..753dfed9 --- /dev/null +++ b/src/lib/cache-invalidation.ts @@ -0,0 +1,139 @@ +import "server-only"; + +import type Redis from "ioredis"; +import { logger } from "@/lib/logger"; +import { redis } from "@/lib/redis"; + +/** + * Cross-process cache signals over Redis pub/sub. + * + * Every process keeps its own in-process fast path, so writing a value to Redis + * is not enough: without a signal, every *other* instance keeps serving its own + * copy until the full TTL elapses. Publishing on a channel lets all instances + * react immediately, which is what makes an invalidation actually take effect. + * + * This is strictly an optimisation. If the subscription never comes up, callers + * still behave correctly and fall back to TTL-based expiry. + */ + +/** Carries a cache key that must be dropped everywhere. */ +export const INVALIDATE_CHANNEL = "cache:invalidate"; + +type Handler = (message: string) => void; + +const handlers = new Map>(); + +// `undefined` = not started yet, `null` = unavailable in this process. +let subscriber: Redis | null | undefined; +let warnedUnavailable = false; + +function canSubscribe(): boolean { + // Never open a connection in the browser, and never during `next build` + // prerendering: a live socket would keep the build process from exiting. + if (typeof window !== "undefined") return false; + if (process.env.NEXT_PHASE === "phase-production-build") return false; + return true; +} + +function dispatch(channel: string, message: string): void { + for (const handler of handlers.get(channel) ?? []) { + try { + handler(message); + } catch (error) { + logger.warn("[cache] signal handler failed", { + channel, + error: String(error), + }); + } + } +} + +function ensureSubscriber(): Redis | null { + if (subscriber !== undefined) return subscriber; + subscriber = null; + if (!redis || !canSubscribe()) return null; + try { + const client: Redis = redis.duplicate(); + client.on("error", (error) => { + logger.warn("[cache] subscriber error", { error: String(error) }); + }); + client.on("message", dispatch); + subscriber = client; + } catch (error) { + logger.warn("[cache] subscriber unavailable", { error: String(error) }); + } + return subscriber; +} + +async function subscribe(channel: string): Promise { + const client = ensureSubscriber(); + if (!client) return; + try { + await client.subscribe(channel); + } catch (error) { + logger.warn("[cache] subscribe failed", { + channel, + error: String(error), + }); + } +} + +const subscribed = new Set(); + +/** + * Run `handler` for every message published on `channel` by any process. + * Returns an unsubscribe function. Safe to call when Redis is absent — the + * handler simply never fires and TTL expiry remains the fallback. + */ +export function onCacheSignal(channel: string, handler: Handler): () => void { + let set = handlers.get(channel); + if (!set) { + set = new Set(); + handlers.set(channel, set); + } + set.add(handler); + if (!subscribed.has(channel)) { + subscribed.add(channel); + void subscribe(channel); + } + return () => { + set.delete(handler); + if (set.size === 0) handlers.delete(channel); + }; +} + +function warnUnavailable(reason: string): void { + if (warnedUnavailable || process.env.NODE_ENV !== "production") return; + warnedUnavailable = true; + logger.warn( + "[cache] Signals cannot reach other instances; entries fall back to TTL expiry.", + { reason }, + ); +} + +/** Publish a message on `channel` for every other process to react to. */ +export async function publishCacheSignal( + channel: string, + message: string, +): Promise { + if (!redis) return; + if (redis.status === "end") { + warnUnavailable("connection closed"); + return; + } + try { + await redis.publish(channel, message); + } catch (error) { + warnUnavailable(String(error)); + } +} + +/** Drop `key` from the in-process cache of every process. */ +export function onCacheInvalidated(handler: Handler): () => void { + return onCacheSignal(INVALIDATE_CHANNEL, handler); +} + +/** Tell every other process to drop `key` from its in-process cache. */ +export async function publishCacheInvalidation(key: string): Promise { + await publishCacheSignal(INVALIDATE_CHANNEL, key); +} diff --git a/src/lib/cache-stats.test.ts b/src/lib/cache-stats.test.ts new file mode 100644 index 00000000..a423ddd5 --- /dev/null +++ b/src/lib/cache-stats.test.ts @@ -0,0 +1,99 @@ +// @ts-nocheck +import { beforeEach, expect, it, vi } from "vitest"; + +const state = vi.hoisted(() => ({ + status: "ready", + get: vi.fn(), + setex: vi.fn(), +})); + +vi.mock("@/lib/redis", () => ({ + redis: { + get: (key: string) => state.get(key), + setex: (key: string, ttl: number, value: string) => + state.setex(key, ttl, value), + get status() { + return state.status; + }, + }, +})); +vi.mock("@/lib/logger", () => ({ + logger: { warn: vi.fn(), error: vi.fn(), info: vi.fn() }, +})); + +import { readCacheStats, recordCacheOutcome } from "./cache-stats"; + +beforeEach(() => { + // The map is process-global on purpose; tests get a clean one. + globalThis.cacheStats?.clear(); + state.status = "ready"; + state.get.mockReset().mockResolvedValue(null); + state.setex.mockReset().mockResolvedValue("OK"); +}); + +it("counts outcomes per key", async () => { + recordCacheOutcome("online_count", "hit"); + recordCacheOutcome("online_count", "hit"); + recordCacheOutcome("online_count", "miss"); + recordCacheOutcome("online_count", "stale"); + recordCacheOutcome("api_staff", "error"); + recordCacheOutcome("api_staff", "evicted"); + + const report = await readCacheStats(); + expect(report.keys.online_count).toEqual({ + hits: 2, + misses: 1, + stale: 1, + errors: 0, + evictions: 0, + }); + expect(report.keys.api_staff.errors).toBe(1); + expect(report.keys.api_staff.evictions).toBe(1); + expect(report.totals.keys).toBe(2); +}); + +it("counts a stale serve as a hit, since the origin was not called", async () => { + recordCacheOutcome("k", "hit"); + recordCacheOutcome("k", "stale"); + recordCacheOutcome("k", "stale"); + const { totals } = await readCacheStats(); + // 1 real hit + 2 stale serves, none of which reached the database. + expect(totals.hits).toBe(1); + expect(totals.stale).toBe(2); + expect(totals.hitRatio).toBe(1); +}); + +it("keeps the tracked key space bounded", async () => { + for (let i = 0; i < 600; i++) recordCacheOutcome(`key-${i}`, "hit"); + const report = await readCacheStats(); + expect(report.totals.keys).toBeLessThanOrEqual(500); +}); + +it("reports Redis as unavailable when it is down, and stays usable", async () => { + state.status = "end"; + recordCacheOutcome("online_count", "hit"); + const report = await readCacheStats(); + expect(report.redisAvailable).toBe(false); + expect(report.shared).toBe(false); + expect(report.keys.online_count.hits).toBe(1); +}); + +it("prefers the shared snapshot from Redis when one exists", async () => { + recordCacheOutcome("local_only", "hit"); + state.get.mockResolvedValue( + JSON.stringify({ shared_key: { hits: 9, misses: 1 } }), + ); + const report = await readCacheStats(); + expect(report.shared).toBe(true); + expect(report.keys.shared_key.hits).toBe(9); + expect(report.totals.hitRatio).toBe(0.9); +}); + +it("survives an unreadable shared snapshot", async () => { + recordCacheOutcome("local_only", "hit"); + state.get.mockRejectedValue(Error("offline")); + const report = await readCacheStats(); + expect(report.shared).toBe(false); + expect(report.redisAvailable).toBe(true); + expect(report.keys.local_only.hits).toBe(1); +}); diff --git a/src/lib/cache-stats.ts b/src/lib/cache-stats.ts new file mode 100644 index 00000000..70466a8b --- /dev/null +++ b/src/lib/cache-stats.ts @@ -0,0 +1,154 @@ +import "server-only"; + +import { logger } from "@/lib/logger"; +import { redis } from "@/lib/redis"; + +/** + * Hit/miss accounting for the in-process/Redis cache. + * + * Without this there is no way to tell a healthy cache from one that is silently + * serving nothing: a wrong Redis URL, a full memory budget or a dead origin all + * look identical from the outside — everything just gets slower. Counters are + * aggregated in-process and flushed to Redis on a timer, so the numbers survive + * a restart and can be read from any instance. + */ + +const REDIS_KEY = "cms:cache-stats:v1:shared"; +const FLUSH_INTERVAL_MS = 30_000; +const WINDOW_MS = 5 * 60 * 1000; + +export interface CacheKeyStats { + hits: number; + misses: number; + /** Served past the TTL from the grace window while refreshing behind it. */ + stale: number; + /** Origin calls that threw. */ + errors: number; + /** Evictions from the in-process budget, i.e. a key we had to recompute. */ + evictions: number; +} + +export type CacheOutcome = "hit" | "miss" | "stale" | "error" | "evicted"; + +const OUTCOME_FIELD = { + hit: "hits", + miss: "misses", + stale: "stale", + error: "errors", + evicted: "evictions", +} as const satisfies Record; + +const globalForStats = globalThis as typeof globalThis & { + cacheStats?: Map; + cacheStatsTimer?: ReturnType; +}; + +// Per key, not global: one hot key thrashing must be visible on its own. +const stats = globalForStats.cacheStats ?? new Map(); +globalForStats.cacheStats = stats; + +// Bound the key space: a high-cardinality cache must not leak one entry per key. +const MAX_TRACKED_KEYS = 500; + +function empty(): CacheKeyStats { + return { hits: 0, misses: 0, stale: 0, errors: 0, evictions: 0 }; +} + +function slot(key: string): CacheKeyStats { + let entry = stats.get(key); + if (entry) return entry; + if (stats.size >= MAX_TRACKED_KEYS) { + // Map order is insertion order and record() re-inserts, so the first key + // is the least recently accounted one. + const oldest = stats.keys().next().value; + if (oldest !== undefined) stats.delete(oldest); + } + entry = empty(); + stats.set(key, entry); + return entry; +} + +/** Record the outcome of one `cached()` call. Never throws, never awaits. */ +export function recordCacheOutcome(key: string, outcome: CacheOutcome): void { + const entry = slot(key); + entry[OUTCOME_FIELD[outcome]]++; + stats.delete(key); + stats.set(key, entry); + schedule(); +} + +function snapshot(): Record { + return Object.fromEntries( + [...stats.entries()].map(([key, value]) => [key, { ...value }]), + ); +} + +async function flush(): Promise { + if (redis?.status !== "ready") return; + try { + await redis.setex( + REDIS_KEY, + Math.ceil(WINDOW_MS / 1000), + JSON.stringify(snapshot()), + ); + } catch { + // Metrics must never interfere with serving. + } +} + +function schedule(): void { + if (globalForStats.cacheStatsTimer) return; + if (process.env.NEXT_PHASE === "phase-production-build") return; + if (process.env.NODE_ENV === "test") return; + globalForStats.cacheStatsTimer = setInterval(() => { + void flush(); + }, FLUSH_INTERVAL_MS); + // Never hold the process open just to report counters. + globalForStats.cacheStatsTimer.unref?.(); +} + +export interface CacheStatsReport { + /** True when the numbers came from Redis, i.e. cover every instance. */ + shared: boolean; + keys: Record; + totals: CacheKeyStats & { hitRatio: number; keys: number }; + /** False when Redis was unreachable, which silently disables shared caching. */ + redisAvailable: boolean; +} + +export async function readCacheStats(): Promise { + let keys: Record = snapshot(); + let shared = false; + if (redis?.status === "ready") { + try { + const raw = await redis.get(REDIS_KEY); + if (raw) { + keys = JSON.parse(raw) as Record; + shared = true; + } + } catch (error) { + logger.warn("[cache] stats unavailable", { error: String(error) }); + } + } + + const totals = empty(); + for (const value of Object.values(keys)) { + totals.hits += value.hits ?? 0; + totals.misses += value.misses ?? 0; + totals.stale += value.stale ?? 0; + totals.errors += value.errors ?? 0; + totals.evictions += value.evictions ?? 0; + } + const decided = totals.hits + totals.misses + totals.stale; + return { + shared, + keys, + totals: { + ...totals, + // A stale serve is a cache hit: the origin was not called for it. + hitRatio: decided === 0 ? 0 : (totals.hits + totals.stale) / decided, + keys: Object.keys(keys).length, + }, + redisAvailable: Boolean(redis && redis.status === "ready"), + }; +} diff --git a/src/lib/cache.test.ts b/src/lib/cache.test.ts index 137a4586..b399c5dc 100644 --- a/src/lib/cache.test.ts +++ b/src/lib/cache.test.ts @@ -55,13 +55,126 @@ describe("cached (memory-only, no Redis)", () => { expect(fn).toHaveBeenCalledTimes(2); }); - it("evicts the oldest entry when the in-memory cache is full", async () => { + it("evicts the least recently used entry when the cache is full", async () => { const fn = vi.fn(async () => 1); - for (let i = 0; i < 600; i++) { - await cached(`bulk-${i}`, 10_000, fn); + // Fill well past the budget (default 2 000 entries). + for (let i = 0; i < 2_100; i++) { + await cached(`bulk-${i}`, 60_000, fn); } - // Key "bulk-0" was evicted (insertion order), so it must recompute. - await cached("bulk-0", 10_000, fn); - expect(fn).toHaveBeenCalledTimes(601); + // "bulk-0" is the least recently used, so it must have been evicted. + await cached("bulk-0", 60_000, fn); + expect(fn).toHaveBeenCalledTimes(2_101); + }); + + it("keeps a hot key alive while colder keys churn through the cache", async () => { + const hot = vi.fn(async () => "hot"); + const cold = vi.fn(async () => "cold"); + const hotKey = `hot-${Math.random()}`; + + expect(await cached(hotKey, 60_000, hot)).toBe("hot"); + + // Every iteration reads the hot key first, then floods the cache with + // fresh one-off keys. Reading must count as using the entry, so the hot + // key survives even though it was inserted first by a long way. + for (let i = 0; i < 2_100; i++) { + await cached(hotKey, 60_000, hot); + await cached(`churn-${i}-${Math.random()}`, 60_000, cold); + } + + expect(hot).toHaveBeenCalledTimes(1); + }); +}); + +describe("cached (stale-while-revalidate)", () => { + it("serves the stale value and refreshes behind it", async () => { + let count = 0; + const fn = vi.fn(async () => ++count); + const key = `swr-${Math.random()}`; + const ttl = 20; + expect(await cached(key, ttl, fn, { staleMs: 10_000 })).toBe(1); + await new Promise((r) => setTimeout(r, 40)); + // The expired entry is still served, so the caller never waits on the + // origin, and the refresh happens behind the response. + expect(await cached(key, ttl, fn, { staleMs: 10_000 })).toBe(1); + await vi.waitFor(() => expect(fn).toHaveBeenCalledTimes(2)); + expect(await cached(key, ttl, fn, { staleMs: 10_000 })).toBe(2); + }); + + it("recomputes synchronously once the grace window has passed", async () => { + let count = 0; + const fn = vi.fn(async () => ++count); + const key = `swr-expiry-${Math.random()}`; + const ttl = 20; + expect(await cached(key, ttl, fn, { staleMs: 20 })).toBe(1); + await new Promise((r) => setTimeout(r, 80)); + expect(await cached(key, ttl, fn, { staleMs: 20 })).toBe(2); + expect(fn).toHaveBeenCalledTimes(2); + }); + + it("keeps serving the stale value when a background refresh fails", async () => { + const key = `swr-fail-${Math.random()}`; + const fn = vi + .fn() + .mockResolvedValueOnce("first") + .mockRejectedValue(new Error("origin down")); + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + await new Promise((r) => setTimeout(r, 40)); + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + await vi.waitFor(() => expect(fn).toHaveBeenCalledTimes(2)); + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + }); + + it("does not cache the result of a refresh invalidated mid-flight", async () => { + const key = `swr-invalidate-${Math.random()}`; + let resolveSlow: (value: string) => void = () => {}; + const slow = vi.fn( + () => + new Promise((resolve) => { + resolveSlow = resolve; + }), + ); + + const pending = cached(key, 10_000, slow, { staleMs: 10_000 }); + // The refresh is in flight and an invalidation lands before it resolves. + invalidateMemory(key); + resolveSlow("computed-before-invalidation"); + await pending; + + // The outdated value must not have been written back, so the next read + // recomputes instead of serving what the invalidation just discarded. + const fresh = vi.fn(async () => "fresh"); + expect(await cached(key, 10_000, fresh)).toBe("fresh"); + expect(await cached(key, 10_000, fresh)).toBe("fresh"); + expect(fresh).toHaveBeenCalledTimes(1); + }); + + it("recovers once the origin comes back after failed background refreshes", async () => { + const key = `swr-recover-${Math.random()}`; + const fn = vi + .fn() + .mockResolvedValueOnce("first") + .mockRejectedValueOnce(new Error("origin down")) + .mockRejectedValueOnce(new Error("origin down")) + .mockResolvedValue("recovered"); + + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + await new Promise((r) => setTimeout(r, 40)); + + // Each stale read kicks off one background refresh; while the origin keeps + // failing the caller keeps getting the last good value. + for (const expectedCalls of [2, 3]) { + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + await vi.waitFor(() => expect(fn).toHaveBeenCalledTimes(expectedCalls)); + // Let the failed refresh settle so the next read starts a new one + // instead of joining the still-registered in-flight promise. + await new Promise((r) => setTimeout(r, 10)); + } + + // Once a refresh succeeds, the fresh value replaces the stale one. + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("first"); + await vi.waitFor(() => expect(fn).toHaveBeenCalledTimes(4)); + await vi.waitFor(async () => + expect(await cached(key, 20, fn, { staleMs: 10_000 })).toBe("recovered"), + ); }); }); diff --git a/src/lib/cache.ts b/src/lib/cache.ts index 2c5434ac..a7bfc8f9 100644 --- a/src/lib/cache.ts +++ b/src/lib/cache.ts @@ -1,46 +1,143 @@ import "server-only"; +import { onCacheInvalidated } from "@/lib/cache-invalidation"; +import { recordCacheOutcome } from "@/lib/cache-stats"; +import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; -type CacheEntry = { data: T; expiresAt: number }; -const memory = new Map>(); +type CacheEntry = { + data: unknown; + /** Fresh until this timestamp. */ + expiresAt: number; + /** Stale values stay servable until this timestamp (stale-while-revalidate). */ + staleUntil: number; +}; -// Cap the in-process map so dynamic keys (leaderboard currencies, article -// slugs, …) can never grow it without bound. Oldest entries are evicted. -const MAX_MEMORY_ENTRIES = 500; +export interface CachedOptions { + /** + * How long past the TTL the last value may still be served while it is + * refreshed in the background. `0` (the default) recomputes synchronously on + * expiry, which is what callers that need an immediately-fresh value want. + * A small grace turns a TTL boundary from a blocking origin call into a + * non-blocking one, so a mass expiry can never stall requests. + * + * The window is stored on the entry, not on the caller, so a key only needs + * *one* call site to opt in: every other reader of that key benefits from the + * grace window too. That matters for hot shared keys like `online_count`, + * which are read from a dozen places — opt in once, centrally. + */ + staleMs?: number; +} + +const memory = new Map(); + +function readPositiveInt(name: string, fallback: number): number { + const raw = (process.env[name] || "").trim(); + if (!raw) return fallback; + const parsed = Number.parseInt(raw, 10); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : fallback; +} + +// In-process budget. Entries are evicted least-recently-*used* (see getMemory), +// so a hot key only leaves when a hotter one takes its place. The default is +// generous because the layer exists to keep origin calls off the hot path. +const MAX_MEMORY_ENTRIES = readPositiveInt("CACHE_MEMORY_MAX_ENTRIES", 2_000); // Single-flight: a key being (re)computed is awaited by concurrent callers // instead of each starting its own `fn()` (cache-stampede protection). const inFlight = new Map>(); +// Bumped by every invalidation. A refresh that started before the bump must not +// write its (now outdated) result back into the cache, or an invalidation +// would be silently undone by a refresh that was already in flight. +const generations = new Map(); +const MAX_GENERATION_ENTRIES = MAX_MEMORY_ENTRIES * 2; + +// Keys whose most recent background refresh failed, so a persistently broken +// origin reports once instead of on every read of a hot key. +const reportedFailures = new Set(); + /** Drop a key from the in-process cache (used when an upstream value changes). */ export function invalidateMemory(key: string): void { + generations.set(key, (generations.get(key) ?? 0) + 1); memory.delete(key); inFlight.delete(key); + reportedFailures.delete(key); + // A generation only matters to a refresh that is still running; once nothing + // is in flight and the value is gone, the counter is dead weight. + if (generations.size > MAX_GENERATION_ENTRIES) { + for (const key of generations.keys()) { + if (!memory.has(key) && !inFlight.has(key)) generations.delete(key); + if (generations.size <= MAX_GENERATION_ENTRIES) break; + } + } } function pruneExpired(now: number): void { for (const [key, entry] of memory) { - if (entry.expiresAt <= now) memory.delete(key); + if (entry.staleUntil <= now) memory.delete(key); } } -function setMemory(key: string, entry: CacheEntry): void { +function setMemory(key: string, entry: CacheEntry): void { if (memory.has(key)) { memory.set(key, entry); return; } + // Reclaim before evicting so a full map of long-lived entries never has to + // drop a live one to make room. if (memory.size >= MAX_MEMORY_ENTRIES) { - // Drop expired entries first, then evict the oldest (insertion order). pruneExpired(Date.now()); if (memory.size >= MAX_MEMORY_ENTRIES) { - const oldest = memory.keys().next().value; - if (oldest !== undefined) memory.delete(oldest); + // Map iteration order is insertion order and getMemory() re-inserts + // on every read, so the first key is the least recently used. + const lru = memory.keys().next().value; + if (lru !== undefined) { + memory.delete(lru); + recordCacheOutcome(lru, "evicted"); + } } } memory.set(key, entry); } +/** Current in-process budget usage, for the admin cache report. */ +export function cacheBudget(): { entries: number; limit: number } { + return { entries: memory.size, limit: MAX_MEMORY_ENTRIES }; +} + +/** + * Read a key and mark it as most recently used. Re-inserting is what makes the + * eviction above true LRU: without it, a key that is read on every single + * request still sits at the front of the insertion order and gets evicted by an + * unrelated burst of dynamic keys, which looks exactly like random clearing. + */ +function getMemory(key: string): CacheEntry | undefined { + const entry = memory.get(key); + if (entry === undefined) return undefined; + memory.delete(key); + memory.set(key, entry); + return entry; +} + +/** `setex` with a little jitter so keys written together do not expire together. */ +function redisTtlSeconds(ttlSec: number): number { + const jitter = Math.min(5, Math.floor(ttlSec * 0.1)); + return ttlSec + Math.floor(Math.random() * (jitter + 1)); +} + +// A wrong REDIS_URL or a Redis outage does not fail visibly: the cache keeps +// answering from memory and every instance quietly stops sharing. Say so once. +let warnedSharedCacheDown = false; +function warnIfSharedCacheDown(): void { + if (redis && redis.status !== "end") return; + if (warnedSharedCacheDown || process.env.NODE_ENV !== "production") return; + warnedSharedCacheDown = true; + logger.error( + "[cache] Shared (Redis) caching is unavailable. Entries are per-process only, so they are dropped on restart and not shared between instances. Check REDIS_URL and that Redis is reachable.", + ); +} + /** * Redis-first cached query with an in-memory fallback. * Use for read-heavy endpoints polled by the browser (online count, etc.). @@ -49,47 +146,103 @@ export async function cached( key: string, ttlMs: number, fn: () => Promise, + options: CachedOptions = {}, ): Promise { - const ttlSec = Math.ceil(ttlMs / 1000); - - // In-memory fast path — served first so repeated reads within a TTL window - // don't each pay a Redis round-trip (Redis is still the shared fallback). + const staleMs = Math.max(0, options.staleMs ?? 0); // `existing &&` short-circuits so Date.now() is never evaluated during // prerender when the map is empty (keeps `next build` prerendering clean). - const existing = memory.get(key); - if (existing && existing.expiresAt > Date.now()) { - return existing.data as T; + const existing = getMemory(key); + if (existing) { + const now = Date.now(); + if (existing.expiresAt > now) { + recordCacheOutcome(key, "hit"); + return existing.data as T; + } + if (existing.staleUntil > now) { + recordCacheOutcome(key, "stale"); + // Serve the last good value and refresh behind it: the caller never + // waits on the origin, and concurrent readers share one refresh. + if (!inFlight.has(key)) { + void refresh(key, ttlMs, staleMs, fn).catch((error: unknown) => { + // A failed refresh keeps the stale entry until its grace window + // runs out, then the next read recomputes synchronously. Never + // let it surface as an unhandled rejection, and never log it per + // read: a broken origin would otherwise flood the log from a + // single hot key. + if (reportedFailures.has(key)) return; + reportedFailures.add(key); + logger.warn("[cache] background refresh failed", { + key, + error: String(error), + }); + }); + } + return existing.data as T; + } + memory.delete(key); } // Single-flight: a concurrent request already recomputing this key. const pending = inFlight.get(key); - if (pending) return (await pending) as T; + if (pending) { + recordCacheOutcome(key, "miss"); + return (await pending) as T; + } + recordCacheOutcome(key, "miss"); + return refresh(key, ttlMs, staleMs, fn); +} + +/** + * Recompute `key` and publish the result. Runs behind a single-flight guard, so + * concurrent misses (and background refreshes) share one origin call. + */ +function refresh( + key: string, + ttlMs: number, + staleMs: number, + fn: () => Promise, +): Promise { + const generation = generations.get(key) ?? 0; const compute = (async (): Promise => { // Redis path (shared across instances). if (redis && redis.status !== "end") { try { - const cached = await redis.get(key); - if (cached !== null && cached !== undefined) { - const data = JSON.parse(cached) as T; - setMemory(key, { data, expiresAt: Date.now() + ttlMs }); + const raw = await redis.get(key); + if (raw !== null && raw !== undefined) { + const data = JSON.parse(raw) as T; + commit(key, data, ttlMs, staleMs, generation); return data; } } catch { /* fall through to fn */ } + } else { + warnIfSharedCacheDown(); } - const data = await fn(); + const data = await fn().catch((error: unknown) => { + recordCacheOutcome(key, "error"); + throw error; + }); - if (redis && redis.status !== "end") { - try { - await redis.setex(key, ttlSec, JSON.stringify(data)); - } catch { - /* non-critical: memory cache still works */ + // An invalidation landed while `fn()` was running: keep the value out of + // the cache so the next read recomputes instead of resurrecting stale data. + if ((generations.get(key) ?? 0) === generation) { + if (redis && redis.status !== "end") { + try { + await redis.setex( + key, + redisTtlSeconds(Math.ceil(ttlMs / 1000)), + JSON.stringify(data), + ); + } catch { + /* non-critical: memory cache still works */ + } } + reportedFailures.delete(key); + commit(key, data, ttlMs, staleMs, generation); } - setMemory(key, { data, expiresAt: Date.now() + ttlMs }); return data; })().finally(() => { @@ -99,3 +252,24 @@ export async function cached( inFlight.set(key, compute); return compute; } + +function commit( + key: string, + data: unknown, + ttlMs: number, + staleMs: number, + generation: number, +): void { + if ((generations.get(key) ?? 0) !== generation) return; + const now = Date.now(); + setMemory(key, { + data, + expiresAt: now + ttlMs, + staleUntil: now + ttlMs + staleMs, + }); +} + +// Another process (the jobs worker, a second instance) invalidating a key must +// drop this process's copy too, otherwise a value stays stale here for the full +// TTL after everyone else already saw the new one. +onCacheInvalidated(invalidateMemory); diff --git a/src/lib/cached-db.test.ts b/src/lib/cached-db.test.ts new file mode 100644 index 00000000..f0168f5f --- /dev/null +++ b/src/lib/cached-db.test.ts @@ -0,0 +1,69 @@ +// @ts-nocheck +import { beforeEach, expect, it, vi } from "vitest"; + +const state = vi.hoisted(() => ({ + deleted: [] as string[], + published: [] as string[], + failDelete: false, + status: "ready", + localKeys: new Set(), +})); + +vi.mock("@/lib/redis", () => ({ + redis: { + get status() { + return state.status; + }, + del: async (key: string) => { + if (state.failDelete) throw Error("redis write failed"); + state.deleted.push(key); + return 1; + }, + publish: async (channel: string, key: string) => { + state.published.push(`${channel}|${key}`); + return 1; + }, + duplicate: () => ({ + on: () => {}, + subscribe: async () => 1, + }), + }, +})); +vi.mock("@/lib/logger", () => ({ + logger: { warn: vi.fn(), error: vi.fn(), info: vi.fn() }, +})); +vi.mock("@/lib/cache-invalidation", () => ({ + publishCacheInvalidation: async (key: string) => { + state.published.push(`broadcast|${key}`); + }, + onCacheInvalidated: () => {}, +})); + +import { invalidateKey } from "./cached-db"; + +beforeEach(() => { + state.deleted = []; + state.published = []; + state.failDelete = false; + state.status = "ready"; +}); + +it("clears memory, Redis and the other instances", async () => { + expect(await invalidateKey("home:news")).toBe(1); + expect(state.deleted).toEqual(["home:news"]); + expect(state.published).toEqual(["broadcast|home:news"]); +}); + +it("still tells the other instances when the Redis delete fails", async () => { + // Returning early here would leave every other process serving the old value + // until its own TTL ran out, while this one already had the new value. + state.failDelete = true; + expect(await invalidateKey("home:news")).toBe(0); + expect(state.published).toEqual(["broadcast|home:news"]); +}); + +it("still tells the other instances when Redis is gone", async () => { + state.status = "end"; + expect(await invalidateKey("home:news")).toBe(0); + expect(state.published).toEqual(["broadcast|home:news"]); +}); diff --git a/src/lib/cached-db.ts b/src/lib/cached-db.ts index 0b25d5f0..25f84478 100644 --- a/src/lib/cached-db.ts +++ b/src/lib/cached-db.ts @@ -1,6 +1,8 @@ import "server-only"; -import { cached, invalidateMemory } from "@/lib/cache"; +import { type CachedOptions, cached, invalidateMemory } from "@/lib/cache"; +import { publishCacheInvalidation } from "@/lib/cache-invalidation"; +import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; const DEFAULT_CACHE_TTL = 60; @@ -14,19 +16,28 @@ export async function cachedQuery( cacheKey: string, queryFn: () => Promise, ttl: number = DEFAULT_CACHE_TTL, + options?: CachedOptions, ): Promise { - return cached(cacheKey, ttl * 1000, queryFn); + return cached(cacheKey, ttl * 1000, queryFn, options); } -/** Invalidate a single cache key immediately (memory + Redis). */ +/** Invalidate a single cache key immediately (memory + Redis + every instance). */ export async function invalidateKey(key: string): Promise { invalidateMemory(key); + let removed = 0; if (redis && redis.status !== "end") { try { - return await redis.del(key); - } catch { - return 0; + removed = await redis.del(key); + } catch (error) { + // A failed delete is not a reason to stay quiet: the durable copy may + // still exist, but the other processes have to drop theirs either way, + // otherwise this instance is the only one that sees the new value. + logger.warn("[cache] redis delete failed", { + key, + error: String(error), + }); } } - return 0; + await publishCacheInvalidation(key); + return removed; } diff --git a/src/lib/imager-cache.test.ts b/src/lib/imager-cache.test.ts index 9e98bf9c..e6011bd4 100644 --- a/src/lib/imager-cache.test.ts +++ b/src/lib/imager-cache.test.ts @@ -1,5 +1,7 @@ // @ts-nocheck -import { beforeEach, describe, expect, it, vi } from "vitest"; +import { mkdir, readdir, rm, writeFile } from "node:fs/promises"; +import { join } from "node:path"; +import { afterAll, beforeEach, describe, expect, it, vi } from "vitest"; import { avatarCacheDir, imagingCacheKey, @@ -10,12 +12,18 @@ import { const TEST_KEY = imagingCacheKey("test-render"); const TEST_DIR = () => `${avatarCacheDir()}/unit`; +// One scratch root for the whole file. A per-test root (pid + Date.now()) left a +// new directory behind for every single test, which is how a hundred+ junk +// directories ended up in the runtime cache tree. +const SCRATCH_ROOT = `${process.cwd()}/storage/imaging/test-unit-${process.pid}`; + beforeEach(() => { // Isolate every run from the runtime cache and from prior runs. - vi.stubEnv( - "IMAGING_CACHE_ROOT", - `${process.cwd()}/storage/imaging/test-unit-${process.pid}-${Date.now()}`, - ); + vi.stubEnv("IMAGING_CACHE_ROOT", SCRATCH_ROOT); +}); + +afterAll(async () => { + await rm(SCRATCH_ROOT, { recursive: true, force: true }); }); describe("imagingCacheKey", () => { @@ -56,3 +64,69 @@ describe("readImagingCache / writeImagingCache", () => { expect(await readImagingCache(TEST_DIR(), key)).not.toBeNull(); }); }); + +describe("pruneImagingCache", () => { + it("brings a directory that is over budget back down to its low-water mark", async () => { + // The budget is read per sweep, so a test can make it small instead of + // writing twenty thousand files. + vi.stubEnv("IMAGING_CACHE_MAX_ENTRIES", "20"); + const dir = `${SCRATCH_ROOT}/prune-cap`; + + // Fill past the budget before the first write of this directory, because + // that first write is the one that runs the sweep. + await mkdir(dir, { recursive: true }); + for (let i = 0; i < 30; i++) { + const name = imagingCacheKey(`figure-${String(i).padStart(3, "0")}`); + await writeFile(join(dir, `${name}.img`), "x"); + await writeFile(join(dir, `${name}.json`), "{}"); + } + + // Every file is fresh, so an age-based sweep would not remove one of + // them and the directory would keep growing. The write triggers a sweep. + await writeImagingCache( + dir, + imagingCacheKey("trigger-sweep"), + new Uint8Array([1]), + "image/png", + ); + + // Back inside the budget: 20 * 0.9 = 18 survivors, plus the record the + // triggering write just added. Records are counted as .img/.json pairs, + // not as individual files. + const remaining = (await readdir(dir)).filter((n) => n.endsWith(".img")); + expect(remaining.length).toBeLessThanOrEqual(19); + expect(remaining.length).toBeGreaterThan(0); + + // Newest survive, oldest go: the prune is least-recently-used, not random. + expect(remaining).toContain(`${imagingCacheKey("figure-029")}.img`); + expect(remaining).not.toContain(`${imagingCacheKey("figure-000")}.img`); + + // A pair is removed together, so no .json is left without its image. + for (const image of remaining) + expect(await readdir(dir)).toContain( + `${image.replace(/\.img$/, "")}.json`, + ); + }); + + it("leaves a directory that is within budget alone", async () => { + vi.stubEnv("IMAGING_CACHE_MAX_ENTRIES", "20"); + const dir = `${SCRATCH_ROOT}/prune-under-cap`; + for (let i = 0; i < 5; i++) { + await writeImagingCache( + dir, + imagingCacheKey(`kept-${i}`), + new Uint8Array([i]), + "image/png", + ); + } + const before = await readdir(dir); + await writeImagingCache( + dir, + imagingCacheKey("another"), + new Uint8Array([1]), + "image/png", + ); + // Only the new record was added; nothing that was already there was dropped. + expect((await readdir(dir)).length).toBe(before.length + 2); + }); +}); diff --git a/src/lib/imager-cache.ts b/src/lib/imager-cache.ts index 6857e5db..bbb3d484 100644 --- a/src/lib/imager-cache.ts +++ b/src/lib/imager-cache.ts @@ -10,7 +10,7 @@ import { unlink, writeFile, } from "node:fs/promises"; -import { join } from "node:path"; +import { extname, join } from "node:path"; /** * Persistent disk cache for imaging renders (avatars, badges). @@ -25,11 +25,30 @@ import { join } from "node:path"; * the directory cannot grow without bound. */ -const MAX_ENTRIES = 20_000; +function maxEntries(): number { + const raw = (process.env.IMAGING_CACHE_MAX_ENTRIES || "").trim(); + const parsed = Number.parseInt(raw, 10); + return Number.isSafeInteger(parsed) && parsed > 0 ? parsed : 20_000; +} + const MAX_AGE_MS = 30 * 24 * 60 * 60 * 1000; +// Pruning reads and stats every file in the directory, so running it on every +// write costs thousands of syscalls per rendered avatar. A write only needs the +// sweep to have happened "recently enough" — the budget is enforced over hours, +// not per request, so the exact moment does not matter. +const PRUNE_INTERVAL_MS = 5 * 60 * 1000; + +// Prune down to this fraction of the budget, so the sweep only runs again after +// a meaningful amount of new renders rather than on the very next write. +const PRUNE_TARGET_RATIO = 0.9; + const IMG_ROOT_DEFAULT = join(process.cwd(), "storage", "imaging"); +// Keyed per directory so throttling the avatars sweep does not also throttle the +// badges one. +const lastPruneAt = new Map(); + function imagingCacheDir(kind: "avatars" | "badges"): string { // Unit tests point this at a scratch root so they never read or pollute // the runtime cache. @@ -98,21 +117,75 @@ export async function writeImagingCache( } } +/** + * Keep the cache inside its budget. + * + * Records are stored as an `.img`/`.json` pair, so the file count is twice the + * render count. Age alone is not enough to bound the directory: a burst of + * popular figures can push it far past the budget while every file is still + * fresh, and then nothing is ever removed. So expired files go first, and if that + * is not enough the oldest remaining files go too, down to a low-water mark so + * the next sweep is not immediately due again. + */ async function pruneImagingCache(directory: string): Promise { try { + const now = Date.now(); + const lastRun = lastPruneAt.get(directory) ?? 0; + if (now - lastRun < PRUNE_INTERVAL_MS) return; + lastPruneAt.set(directory, now); + const names = await readdir(directory); - if (names.length <= MAX_ENTRIES) return; - const cutoffMs = Date.now() - MAX_AGE_MS; + // Group the two halves of each record so a pair is never left orphaned. + const records = new Map(); for (const name of names) { + const base = name.slice(0, name.length - extname(name).length); + const files = records.get(base); + if (files) files.push(name); + else records.set(base, [name]); + } + const budget = maxEntries(); + if (records.size <= budget) return; + + const cutoffMs = now - MAX_AGE_MS; + const survivors: { files: string[]; mtimeMs: number }[] = []; + let live = 0; + for (const files of records.values()) { try { - const file = join(directory, name); - const info = await stat(file); - if (info.mtimeMs < cutoffMs) await unlink(file); + const info = await stat(join(directory, files[0] as string)); + if (info.mtimeMs < cutoffMs) { + await removeRecord(directory, files); + continue; + } + survivors.push({ files, mtimeMs: info.mtimeMs }); + live++; } catch { - // Skip entries that disappeared between listing and unlink. + // Skip records that disappeared between listing and stat. + } + } + + const target = Math.floor(budget * PRUNE_TARGET_RATIO); + if (live <= target) return; + survivors.sort((a, b) => a.mtimeMs - b.mtimeMs); + for (const record of survivors) { + if (live <= target) break; + try { + await removeRecord(directory, record.files); + live--; + } catch { + // Already gone. } } } catch { // Nothing to prune or directory missing. } } + +async function removeRecord(directory: string, files: string[]): Promise { + for (const name of files) { + try { + await unlink(join(directory, name)); + } catch { + // Already gone. + } + } +} diff --git a/src/lib/redis-cache.ts b/src/lib/redis-cache.ts index 9c4ae681..4a601446 100644 --- a/src/lib/redis-cache.ts +++ b/src/lib/redis-cache.ts @@ -1,6 +1,6 @@ import "server-only"; -import { cached } from "@/lib/cache"; +import { type CachedOptions, cached } from "@/lib/cache"; /** * Cache the result of a fetch function, memory-first with a Redis fallback @@ -10,8 +10,9 @@ export async function redisCache( key: string, ttlSeconds: number, fetch: () => Promise, + options?: CachedOptions, ): Promise { - return cached(key, ttlSeconds * 1000, fetch); + return cached(key, ttlSeconds * 1000, fetch, options); } /** diff --git a/src/lib/runtime-avatar-route.test.ts b/src/lib/runtime-avatar-route.test.ts index 1d80c484..8d0cb680 100644 --- a/src/lib/runtime-avatar-route.test.ts +++ b/src/lib/runtime-avatar-route.test.ts @@ -1,21 +1,27 @@ // @ts-nocheck +import { rm } from "node:fs/promises"; import { NextRequest } from "next/server"; -import { afterEach, beforeEach, expect, it, vi } from "vitest"; +import { afterAll, afterEach, beforeEach, expect, it, vi } from "vitest"; import { GET } from "../app/api/imaging/avatar/route"; +// One scratch root for the whole file, removed afterwards: a per-test root left a +// new directory behind for every test, polluting the runtime cache tree. +const SCRATCH_ROOT = `${process.cwd()}/storage/imaging/test-route-${process.pid}`; + beforeEach(() => { - // Fresh, isolated cache root per test run so the disk cache can never leak - // state between runs or into the runtime cache. - vi.stubEnv( - "IMAGING_CACHE_ROOT", - `${process.cwd()}/storage/imaging/test-route-${process.pid}-${Date.now()}`, - ); + // Isolated cache root so the disk cache can never leak state between runs or + // into the runtime cache. + vi.stubEnv("IMAGING_CACHE_ROOT", SCRATCH_ROOT); }); afterEach(() => { vi.unstubAllEnvs(); vi.unstubAllGlobals(); }); + +afterAll(async () => { + await rm(SCRATCH_ROOT, { recursive: true, force: true }); +}); it("forwards multicolor figures and preserves configured query options at runtime", async () => { const figure = "hr-11782-40-40.hd-180-7-14.ch-11592-66.lg-10726-79-1408"; const fetcher = vi.fn( diff --git a/src/lib/services/cache-warmup.test.ts b/src/lib/services/cache-warmup.test.ts new file mode 100644 index 00000000..007b0140 --- /dev/null +++ b/src/lib/services/cache-warmup.test.ts @@ -0,0 +1,104 @@ +// @ts-nocheck +import { beforeEach, expect, it, vi } from "vitest"; + +const state = vi.hoisted(() => ({ + // These stand in for the real cache, which runs the query on a miss, so a + // failing database really does propagate into the warm-up. + cached: vi.fn( + async (_key: string, _ttl: number, query: () => Promise) => + query(), + ), + redisCache: vi.fn( + async (_key: string, _ttl: number, query: () => Promise) => + query(), + ), + cacheNews: vi.fn( + async (_key: string, _ttl: number, query: () => Promise) => + query(), + ), + warn: vi.fn(), + // Set by a test to make every origin query reject. + failEverything: false, +})); + +vi.mock("@/lib/logger", () => ({ + logger: { warn: state.warn, error: vi.fn(), info: vi.fn() }, +})); +vi.mock("@/lib/cache", () => ({ cached: state.cached })); +vi.mock("@/lib/redis-cache", () => ({ + redisCache: state.redisCache, + apiCacheKey: (name: string) => `api:${name}`, +})); +vi.mock("@/lib/services/news-cache", () => ({ cacheNews: state.cacheNews })); +vi.mock("@/lib/services/site-settings", () => ({ + siteSettings: { get: vi.fn(async () => "7") }, +})); +vi.mock("@/lib/db", () => { + // A chainable stand-in for a drizzle query that always resolves to a row. + const row = { total: 1, username: "a", look: "l", rank: 7, motto: "m" }; + // Chainable and awaitable at once, the way a drizzle query builder is. + const chain: any = new Proxy(() => {}, { + get: (_target, prop) => { + if (prop === Symbol.toStringTag) return "Query"; + if (prop === "then") + return (resolve: any, reject: any) => + (state.failEverything + ? Promise.reject(Error("database is down")) + : Promise.resolve([row]) + ).then(resolve, reject); + return () => chain; + }, + }); + return { + db: { select: () => chain }, + User: { online: "online", username: "username", rank: "rank" }, + Rooms: {}, + CameraWeb: {}, + WebsiteTeams: { id: "id", hiddenRank: "hiddenRank" }, + }; +}); + +import { warmPublicCaches } from "./cache-warmup"; + +beforeEach(() => { + state.failEverything = false; + for (const mock of [ + state.cached, + state.redisCache, + state.cacheNews, + state.warn, + ]) + mock.mockClear(); +}); + +it("primes the keys the public pages actually read", async () => { + await warmPublicCaches(); + const keys = state.cached.mock.calls.map(([key]) => key); + expect(keys).toContain("online_count"); + expect(keys).toContain("total_users"); + expect(keys).toContain("total_rooms"); + expect(keys).toContain("total_photos"); + expect(state.redisCache.mock.calls.map(([key]) => key)).toEqual( + expect.arrayContaining(["api:staff", "api:teams"]), + ); + expect(state.cacheNews).toHaveBeenCalled(); +}); + +it("gives every primed key a stale window", async () => { + // Warming without a grace window would leave the first visitor after a deploy + // to still pay for a miss, which is the thing being fixed. + await warmPublicCaches(); + for (const call of [ + ...state.cached.mock.calls, + ...state.redisCache.mock.calls, + ]) + expect(call[3]?.staleMs).toBeGreaterThan(0); +}); + +it("survives a database that is completely unavailable", async () => { + // register() calls this without awaiting: it must never reject, or an + // unhandled rejection would take the process down on every boot. + state.failEverything = true; + await expect(warmPublicCaches()).resolves.toBeUndefined(); + expect(state.warn).toHaveBeenCalled(); +}); diff --git a/src/lib/services/cache-warmup.ts b/src/lib/services/cache-warmup.ts new file mode 100644 index 00000000..4ce0e0a5 --- /dev/null +++ b/src/lib/services/cache-warmup.ts @@ -0,0 +1,135 @@ +import "server-only"; + +import { asc, count, desc, eq, gte } from "drizzle-orm"; +import { cached } from "@/lib/cache"; +import { CameraWeb, db, Rooms, User, WebsiteTeams } from "@/lib/db"; +import { logger } from "@/lib/logger"; +import { apiCacheKey, redisCache } from "@/lib/redis-cache"; +import { cacheNews } from "@/lib/services/news-cache"; +import { siteSettings } from "@/lib/services/site-settings"; + +/** + * Prime the hot public cache keys right after boot. + * + * Every deploy starts with an empty in-process cache, and any key that also + * expired while the old process was down has to be recomputed from the database. + * With no warm-up, the first visitor after each deploy pays for a burst of + * simultaneous misses; with it, the site is already warm before traffic arrives. + * + * Warming works by cache *key*, not by call site: these use exactly the keys the + * routes and pages read, so priming an entry serves every reader of it. Failures + * are logged and ignored — a warm-up that cannot reach the database must never + * stop the server from serving. + */ + +const COUNTER_TTL_MS = 300_000; +const ONLINE_TTL_MS = 10_000; + +async function warm(name: string, load: () => Promise): Promise { + try { + await load(); + } catch (error) { + logger.warn("[cache] warm-up failed", { + key: name, + error: String(error), + }); + } +} + +export async function warmPublicCaches(): Promise { + // Small delays between groups: a boot-time burst of COUNT(*) queries against + // a database that is still opening connections helps nobody. + await warm("online_count", () => + cached("online_count", ONLINE_TTL_MS, countOnline, { staleMs: 15_000 }), + ); + await warm("total_users", () => + cached("total_users", COUNTER_TTL_MS, countUsers, { + staleMs: COUNTER_TTL_MS, + }), + ); + await warm("total_rooms", () => + cached("total_rooms", COUNTER_TTL_MS, countRooms, { + staleMs: COUNTER_TTL_MS, + }), + ); + await warm("total_photos", () => + cached("total_photos", COUNTER_TTL_MS, countPhotos, { + staleMs: COUNTER_TTL_MS, + }), + ); + await warm("online_users", () => + cached("online_users", ONLINE_TTL_MS, listOnlineUsers, { + staleMs: 15_000, + }), + ); + await warm("api:staff", () => + redisCache(apiCacheKey("staff"), 300, loadStaff, { staleMs: 300 }), + ); + await warm("api:teams", () => + redisCache(apiCacheKey("teams"), 300, loadTeams, { staleMs: 600 }), + ); + await warm("news:home", () => + cacheNews(apiCacheKey("home"), 15_000, async () => ({ primed: true })), + ); +} + +async function countOnline(): Promise { + const [row] = await db + .select({ total: count() }) + .from(User) + .where(eq(User.online, "1")); + return row?.total ?? 0; +} + +async function countUsers(): Promise { + const [row] = await db.select({ total: count() }).from(User); + return row?.total ?? 0; +} + +async function countRooms(): Promise { + const [row] = await db.select({ total: count() }).from(Rooms); + return row?.total ?? 0; +} + +async function countPhotos(): Promise { + const [row] = await db.select({ total: count() }).from(CameraWeb); + return row?.total ?? 0; +} + +async function listOnlineUsers() { + return db + .select({ username: User.username, look: User.look }) + .from(User) + .where(eq(User.online, "1")) + .limit(100); +} + +async function loadStaff() { + const minStaffRank = + Number(await siteSettings.get("min_staff_rank", "7")) || 7; + return db + .select({ + username: User.username, + look: User.look, + rank: User.rank, + motto: User.motto, + }) + .from(User) + .where(gte(User.rank, minStaffRank)) + .orderBy(desc(User.rank), asc(User.username)) + .limit(100); +} + +async function loadTeams() { + return db + .select({ + id: WebsiteTeams.id, + rankName: WebsiteTeams.rankName, + badge: WebsiteTeams.badge, + jobDescription: WebsiteTeams.jobDescription, + staffColor: WebsiteTeams.staffColor, + }) + .from(WebsiteTeams) + .where(eq(WebsiteTeams.hiddenRank, false)) + .orderBy(asc(WebsiteTeams.id)); +} diff --git a/src/lib/services/news-cache.test.ts b/src/lib/services/news-cache.test.ts index b8e104a8..4a000620 100644 --- a/src/lib/services/news-cache.test.ts +++ b/src/lib/services/news-cache.test.ts @@ -6,10 +6,33 @@ const state = vi.hoisted(() => ({ values: new Map(), get: vi.fn(), set: vi.fn(), + publish: vi.fn(), status: "ready", + // Set by the subscriber mock so a test can simulate another process + // publishing a new revision. + deliver: null as ((channel: string, message: string) => void) | null, +})); + +vi.mock("@/lib/redis", () => ({ + redis: { + get: (key: string) => state.get(key), + set: (key: string, value: string) => state.set(key, value), + publish: (channel: string, message: string) => + state.publish(channel, message), + get status() { + return state.status; + }, + duplicate: () => ({ + on: (event: string, handler: (...args: unknown[]) => void) => { + if (event === "message") state.deliver = handler; + }, + subscribe: async () => 1, + }), + }, +})); +vi.mock("@/lib/logger", () => ({ + logger: { error: vi.fn(), warn: vi.fn() }, })); -vi.mock("@/lib/redis", () => ({ redis: state })); -vi.mock("@/lib/logger", () => ({ logger: { error: vi.fn() } })); vi.mock("@/lib/cache", () => ({ cached: async (key: string, _ttl: number, query: () => Promise) => { if (state.values.has(key)) return state.values.get(key); @@ -35,7 +58,9 @@ beforeEach(() => { .mockImplementation(async (_key: string, value: string) => { state.revision = value; }); + state.publish.mockReset().mockResolvedValue(1); }); + it("invalidates a previously cached public list", async () => { const query = vi .fn() @@ -45,6 +70,36 @@ it("invalidates a previously cached public list", async () => { await invalidateNewsCache(); expect(await cacheNews("list", 60000, query)).toEqual(["new article"]); }); + +it("reads the revision at most once for a burst of reads", async () => { + // The revision used to be fetched from Redis on every single call, which put + // a round-trip in front of the very fast path the cache exists to provide. + const before = state.get.mock.calls.length; + await cacheNews("burst-a", 60000, async () => ["a"]); + await cacheNews("burst-b", 60000, async () => ["b"]); + await cacheNews("burst-c", 60000, async () => ["c"]); + expect(state.get.mock.calls.length - before).toBeLessThanOrEqual(1); +}); + +it("picks up a revision published by another process", async () => { + expect(await cacheNews("list", 60000, async () => ["old"])).toEqual(["old"]); + + // Another process rotates the revision; the local copy must be dropped so the + // next read re-reads it instead of serving the previous namespace. + state.revision = "rev-from-worker"; + state.deliver?.("cache:news-revision", "rev-from-worker"); + + expect(await cacheNews("list", 60000, async () => ["new"])).toEqual(["new"]); +}); + +it("tells other processes when the revision rotates", async () => { + await invalidateNewsCache(); + expect(state.publish).toHaveBeenCalledWith( + "cache:news-revision", + expect.any(String), + ); +}); + it("does not let a stale in-flight read replace a newer revision", async () => { let finish!: (value: string[]) => void; const old = cacheNews( @@ -64,6 +119,7 @@ it("does not let a stale in-flight read replace a newer revision", async () => { "new", ]); }); + it("reads fresh data when Redis is unavailable", async () => { state.get.mockRejectedValue(Error("offline")); expect(await cacheNews("list", 60000, async () => ["fresh"])).toEqual([ @@ -71,10 +127,20 @@ it("reads fresh data when Redis is unavailable", async () => { ]); }); +it("still caches in-process when Redis is unavailable", async () => { + // A Redis outage must not turn every public news read into a database query. + state.status = "end"; + const query = vi.fn(async () => ["fresh"]); + expect(await cacheNews("list", 60000, query)).toEqual(["fresh"]); + expect(await cacheNews("list", 60000, query)).toEqual(["fresh"]); + expect(query).toHaveBeenCalledTimes(1); +}); + it("keeps failed durable cache invalidations retryable", async () => { state.set.mockRejectedValue(Error("offline")); await expect(refreshNewsCacheForDelivery()).rejects.toThrow("offline"); }); + it("does not acknowledge an ended Redis connection as refreshed", async () => { state.status = "end"; await expect(refreshNewsCacheForDelivery()).rejects.toThrow(); diff --git a/src/lib/services/news-cache.ts b/src/lib/services/news-cache.ts index a86d69f1..d67da6e8 100644 --- a/src/lib/services/news-cache.ts +++ b/src/lib/services/news-cache.ts @@ -1,28 +1,83 @@ import "server-only"; + import { randomUUID } from "node:crypto"; import { cached } from "@/lib/cache"; +import { onCacheSignal, publishCacheSignal } from "@/lib/cache-invalidation"; import { logger } from "@/lib/logger"; import { redis } from "@/lib/redis"; const REVISION_KEY = "cms:news:revision"; + +/** News entries are namespaced by a revision, so a new revision drops them all. */ +const REVISION_CHANNEL = "cache:news-revision"; + +// The revision is read from Redis at most this often per process. Reading it on +// every call (as this used to) put a Redis round-trip in front of the in-process +// fast path, which is exactly what the cache exists to avoid. A published article +// still shows up immediately: the invalidation signal clears this copy, and the +// window is only the fallback for when pub/sub is unavailable. +const REVISION_TTL_MS = 1_000; + +let revision: { value: string; expiresAt: number } | null = null; +let inFlightRevision: Promise | null = null; + +async function readRevision(): Promise { + const now = Date.now(); + if (revision && revision.expiresAt > now) return revision.value; + inFlightRevision ??= (async () => { + try { + return (await redis?.get(REVISION_KEY)) ?? "0"; + } catch { + return "0"; + } finally { + inFlightRevision = null; + } + })(); + const value = await inFlightRevision; + revision = { value, expiresAt: Date.now() + REVISION_TTL_MS }; + return value; +} + +function setRevision(value: string): void { + revision = { value, expiresAt: Date.now() + REVISION_TTL_MS }; +} + +// Another process published new news: forget the revision so the next read picks +// up the new one instead of serving entries under the old namespace. +onCacheSignal(REVISION_CHANNEL, () => { + revision = null; +}); + export async function cacheNews( key: string, ttlMs: number, query: () => Promise, ): Promise { - if (!redis || redis.status === "end") return query(); - let revision: string; - try { - revision = (await redis.get(REVISION_KEY)) ?? "0"; - } catch { - return query(); - } - return cached(`news:${revision}:${key}`, ttlMs, query); + // Every caller of cacheNews is a public news read, so the grace window lives + // here rather than at each call site. Publishing an article rotates the + // revision, which drops these entries immediately; the window only matters + // when that signal cannot be delivered. + const options = { staleMs: ttlMs }; + // Without Redis the revision cannot be shared, so fall back to a fixed one: + // the in-process cache still works, which is far better than hitting the + // database on every request for the whole duration of a Redis outage. + if (!redis || redis.status === "end") + return cached(`news:0:${key}`, ttlMs, query, options); + return cached(`news:${await readRevision()}:${key}`, ttlMs, query, options); } + +async function rotateRevision(): Promise { + const next = randomUUID(); + if (redis) await redis.set(REVISION_KEY, next); + setRevision(next); + await publishCacheSignal(REVISION_CHANNEL, next); + return next; +} + export async function invalidateNewsCache(): Promise { if (!redis || redis.status === "end") return; try { - await redis.set(REVISION_KEY, randomUUID()); + await rotateRevision(); } catch (error) { logger.error("News saved but public cache invalidation failed", { module: "news", @@ -36,5 +91,5 @@ export async function refreshNewsCacheForDelivery(): Promise { if (!redis) return; if (redis.status === "end") throw new Error("News cache connection is closed"); - await redis.set(REVISION_KEY, randomUUID()); + await rotateRevision(); }