fix(cache): true LRU, stale-while-revalidate and cross-process invalidation
Gitea Actions Runner Test / test-job (push) Successful in 0s
CI / check (push) Successful in 32s
CI / tests-integration (push) Successful in 1m38s
CI / tests-unit (push) Successful in 1m42s
CI / tests-ui (push) Successful in 2m33s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m43s
Gitea Actions Runner Test / test-job (push) Successful in 0s
CI / check (push) Successful in 32s
CI / tests-integration (push) Successful in 1m38s
CI / tests-unit (push) Successful in 1m42s
CI / tests-ui (push) Successful in 2m33s
CI / preflight (push) Skipped
CI / deploy (push) Successful in 2m43s
The in-process cache was a FIFO of 500 entries that was never touched on a read, so a key polled on every request could be evicted by an unrelated burst of dynamic keys. That looked exactly like the cache being cleared at random, and it is what made the site fall back to the database unpredictably. - Evict least-recently-used instead, and raise the default budget to 2000 (CACHE_MEMORY_MAX_ENTRIES). Reading a key now marks it as used, so a hot key only leaves when a hotter one takes its place. - Add opt-in stale-while-revalidate (CachedOptions.staleMs). The grace window lives on the entry, so one call site opting in protects every reader of that key. A failed background refresh keeps serving the last good value instead of falling through to the origin, and is reported once rather than per read. - Invalidate across processes. invalidateKey() now clears memory, deletes the Redis key and publishes a signal, so a value written by one process is no longer served stale by the others for the rest of its TTL. A failed Redis delete no longer skips the broadcast. - Guard against a refresh that started before an invalidation writing its outdated result back into the cache. - Read the news revision at most once a second per process instead of on every call, with a pub/sub signal to drop the local copy when it rotates. A Redis outage now degrades to the in-process cache rather than to no cache at all. - Warm the hot public keys on boot, so the first visitors after a deploy do not each pay for a miss. - Count hits, misses, stale serves, errors and evictions per key, exposed at GET /api/admin/devops/cache. Without it a wrong REDIS_URL, a full budget and a dead origin all look identical from the outside. - Enforce the imaging cache budget for real: records are .img/.json pairs, so the old cap counted files and never removed anything while entries were fresh. Sweeps are throttled per directory and prune to a low-water mark. - Cap the JWT version map, and stop per-test scratch roots from littering the runtime imaging cache. Public read-only endpoints get grace windows; admin, account and auth data deliberately stays fresh. Redis TTLs get a little jitter so keys written together no longer expire together. 3209 tests pass. next build could not be verified on this host: the optimized build is OOM-killed before prerender, so this has not run in a real Next runtime yet.
This commit is contained in:
1 parent
f490fcc9da
commit
203399aab7
43 files changed
+1619
-94
No files matched your search
@@ -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
|
||||
|
||||
@@ -56,7 +56,10 @@ type Row = { username: string; look: string; value: number };
|
||||
|
||||
async function loadCreditsRows(): Promise<Row[]> {
|
||||
try {
|
||||
const users = await cached("lb_credits", 60_000, () =>
|
||||
const users = await cached(
|
||||
"lb_credits",
|
||||
60_000,
|
||||
() =>
|
||||
db
|
||||
.select({
|
||||
username: User.username,
|
||||
@@ -66,6 +69,7 @@ async function loadCreditsRows(): Promise<Row[]> {
|
||||
.from(User)
|
||||
.orderBy(desc(User.credits))
|
||||
.limit(20),
|
||||
{ staleMs: 120000 },
|
||||
);
|
||||
return users.map((u) => ({
|
||||
username: u.username,
|
||||
@@ -79,7 +83,10 @@ async function loadCreditsRows(): Promise<Row[]> {
|
||||
|
||||
async function loadCurrencyRows(type: number): Promise<Row[]> {
|
||||
try {
|
||||
return await cached(`lb_currency_${type}`, 60_000, () =>
|
||||
return await cached(
|
||||
`lb_currency_${type}`,
|
||||
60_000,
|
||||
() =>
|
||||
db
|
||||
.select({
|
||||
username: User.username,
|
||||
@@ -91,6 +98,7 @@ async function loadCurrencyRows(type: number): Promise<Row[]> {
|
||||
.where(eq(UsersCurrency.type, type))
|
||||
.orderBy(desc(UsersCurrency.amount))
|
||||
.limit(20),
|
||||
{ staleMs: 120000 },
|
||||
);
|
||||
} catch {
|
||||
return [];
|
||||
@@ -102,7 +110,10 @@ async function loadSettingsRows(
|
||||
): Promise<Row[]> {
|
||||
try {
|
||||
const column = UsersSettings[field];
|
||||
return await cached(`lb_settings_${field}`, 60_000, () =>
|
||||
return await cached(
|
||||
`lb_settings_${field}`,
|
||||
60_000,
|
||||
() =>
|
||||
db
|
||||
.select({
|
||||
username: User.username,
|
||||
@@ -113,6 +124,7 @@ async function loadSettingsRows(
|
||||
.innerJoin(User, eq(UsersSettings.userId, User.id))
|
||||
.orderBy(desc(column))
|
||||
.limit(20),
|
||||
{ staleMs: 120000 },
|
||||
);
|
||||
} catch {
|
||||
return [];
|
||||
|
||||
+25
-5
@@ -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, () =>
|
||||
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, () =>
|
||||
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, () =>
|
||||
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, () =>
|
||||
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, () =>
|
||||
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")),
|
||||
]);
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -40,7 +40,10 @@ export default async function RankingsPage() {
|
||||
const t = await getTranslations("pages.rankings");
|
||||
let users: TopUser[] | null = [];
|
||||
try {
|
||||
users = await cached("rankings:top", 60_000, async () =>
|
||||
users = await cached(
|
||||
"rankings:top",
|
||||
60_000,
|
||||
async () =>
|
||||
db
|
||||
.select({
|
||||
username: User.username,
|
||||
@@ -52,6 +55,7 @@ export default async function RankingsPage() {
|
||||
.from(User)
|
||||
.orderBy(desc(User.credits))
|
||||
.limit(12),
|
||||
{ staleMs: 120000 },
|
||||
);
|
||||
} catch (error) {
|
||||
users = publicReadFailure("rankings")(error);
|
||||
|
||||
@@ -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, () =>
|
||||
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, () =>
|
||||
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(() => []),
|
||||
]);
|
||||
|
||||
|
||||
+15
@@ -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() });
|
||||
});
|
||||
@@ -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).
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -50,6 +50,7 @@ export async function GET(req: Request) {
|
||||
},
|
||||
});
|
||||
},
|
||||
{ staleMs: 120_000 },
|
||||
);
|
||||
|
||||
return apiJson(data);
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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 () => {
|
||||
// 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 });
|
||||
|
||||
@@ -13,13 +13,19 @@ import { rcon } from "@/lib/services/rcon";
|
||||
|
||||
async function fetchOnlineCount(): Promise<number> {
|
||||
try {
|
||||
return await cached("online_count", 10_000, async () => {
|
||||
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<number> {
|
||||
// degradation on failure so the stream never errors out.
|
||||
async function fetchEmulatorStatus(): Promise<boolean> {
|
||||
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;
|
||||
|
||||
@@ -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 () =>
|
||||
// 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 });
|
||||
|
||||
@@ -44,6 +44,7 @@ export async function GET(req: Request) {
|
||||
},
|
||||
});
|
||||
},
|
||||
{ staleMs: 120_000 },
|
||||
);
|
||||
|
||||
return apiJson(data);
|
||||
|
||||
@@ -38,6 +38,7 @@ export async function GET(_req: Request) {
|
||||
|
||||
return { username: user.username, look: user.look };
|
||||
},
|
||||
{ staleMs: 30_000 },
|
||||
);
|
||||
|
||||
return apiJson({ dj });
|
||||
|
||||
@@ -48,6 +48,7 @@ export async function GET(_req: Request) {
|
||||
};
|
||||
});
|
||||
},
|
||||
{ staleMs: 120_000 },
|
||||
);
|
||||
|
||||
return apiJson({ data });
|
||||
|
||||
@@ -58,6 +58,7 @@ export async function GET(_req: Request) {
|
||||
};
|
||||
});
|
||||
},
|
||||
{ staleMs: 10_000 },
|
||||
);
|
||||
|
||||
return apiJson({ shouts: data });
|
||||
|
||||
@@ -26,6 +26,7 @@ export async function GET(_req: Request) {
|
||||
asc(WebsiteShopCategories.name),
|
||||
),
|
||||
),
|
||||
{ staleMs: 120_000 },
|
||||
);
|
||||
|
||||
return apiJson({ data });
|
||||
|
||||
@@ -71,6 +71,7 @@ export async function GET(req: Request) {
|
||||
},
|
||||
});
|
||||
},
|
||||
{ staleMs: 120_000 },
|
||||
);
|
||||
|
||||
return apiJson(data);
|
||||
|
||||
@@ -12,7 +12,10 @@ export async function GET(_req: Request) {
|
||||
const minStaffRank =
|
||||
Number(await siteSettings.get("min_staff_rank", "7")) || 7;
|
||||
|
||||
const staff = await redisCache(apiCacheKey("staff"), 300, () =>
|
||||
const staff = await redisCache(
|
||||
apiCacheKey("staff"),
|
||||
300,
|
||||
() =>
|
||||
db
|
||||
.select({
|
||||
username: User.username,
|
||||
@@ -24,6 +27,9 @@ export async function GET(_req: Request) {
|
||||
.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 });
|
||||
|
||||
@@ -8,7 +8,10 @@ import { apiCacheKey, cacheSafe, redisCache } from "@/lib/redis-cache";
|
||||
|
||||
export async function GET(_req: Request) {
|
||||
try {
|
||||
const rows = await redisCache(apiCacheKey("teams"), 300, async () =>
|
||||
const rows = await redisCache(
|
||||
apiCacheKey("teams"),
|
||||
300,
|
||||
async () =>
|
||||
cacheSafe(
|
||||
await db
|
||||
.select({
|
||||
@@ -22,6 +25,7 @@ export async function GET(_req: Request) {
|
||||
.where(eq(WebsiteTeams.hiddenRank, false))
|
||||
.orderBy(asc(WebsiteTeams.id)),
|
||||
),
|
||||
{ staleMs: 600_000 },
|
||||
);
|
||||
|
||||
// apiJson serialises BigInt ids → string automatically.
|
||||
|
||||
@@ -76,6 +76,7 @@ export async function GET(
|
||||
|
||||
return cacheSafe({ data: { ...value, category } });
|
||||
},
|
||||
{ staleMs: 300_000 },
|
||||
);
|
||||
|
||||
if (data === null) {
|
||||
|
||||
@@ -26,6 +26,7 @@ export async function GET(_req: Request) {
|
||||
asc(WebsiteRareValueCategories.name),
|
||||
),
|
||||
),
|
||||
{ staleMs: 600_000 },
|
||||
);
|
||||
|
||||
return apiJson({ data });
|
||||
|
||||
@@ -61,6 +61,7 @@ export async function GET(req: Request) {
|
||||
},
|
||||
});
|
||||
},
|
||||
{ staleMs: 300_000 },
|
||||
);
|
||||
|
||||
return apiJson(data);
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
@@ -8,6 +8,24 @@ const MEMORY_TTL_MS = 60_000;
|
||||
const REDIS_TTL_SEC = 60;
|
||||
const memory = new Map<number, { version: number; expiresAt: number }>();
|
||||
|
||||
// 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));
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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<string, Set<Handler>>();
|
||||
|
||||
// `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<void> {
|
||||
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<string>();
|
||||
|
||||
/**
|
||||
* 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<void> {
|
||||
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<void> {
|
||||
await publishCacheSignal(INVALIDATE_CHANNEL, key);
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
@@ -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<CacheOutcome, keyof CacheKeyStats>;
|
||||
|
||||
const globalForStats = globalThis as typeof globalThis & {
|
||||
cacheStats?: Map<string, CacheKeyStats>;
|
||||
cacheStatsTimer?: ReturnType<typeof setInterval>;
|
||||
};
|
||||
|
||||
// Per key, not global: one hot key thrashing must be visible on its own.
|
||||
const stats = globalForStats.cacheStats ?? new Map<string, CacheKeyStats>();
|
||||
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<string, CacheKeyStats> {
|
||||
return Object.fromEntries(
|
||||
[...stats.entries()].map(([key, value]) => [key, { ...value }]),
|
||||
);
|
||||
}
|
||||
|
||||
async function flush(): Promise<void> {
|
||||
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<string, CacheKeyStats>;
|
||||
totals: CacheKeyStats & { hitRatio: number; keys: number };
|
||||
/** False when Redis was unreachable, which silently disables shared caching. */
|
||||
redisAvailable: boolean;
|
||||
}
|
||||
|
||||
export async function readCacheStats(): Promise<CacheStatsReport> {
|
||||
let keys: Record<string, CacheKeyStats> = snapshot();
|
||||
let shared = false;
|
||||
if (redis?.status === "ready") {
|
||||
try {
|
||||
const raw = await redis.get(REDIS_KEY);
|
||||
if (raw) {
|
||||
keys = JSON.parse(raw) as Record<string, CacheKeyStats>;
|
||||
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"),
|
||||
};
|
||||
}
|
||||
+119
-6
@@ -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<string>((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"),
|
||||
);
|
||||
});
|
||||
});
|
||||
+198
-24
@@ -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<T> = { data: T; expiresAt: number };
|
||||
const memory = new Map<string, CacheEntry<unknown>>();
|
||||
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<string, CacheEntry>();
|
||||
|
||||
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<string, Promise<unknown>>();
|
||||
|
||||
// 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<string, number>();
|
||||
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<string>();
|
||||
|
||||
/** 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<T>(key: string, entry: CacheEntry<T>): 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<T>(
|
||||
key: string,
|
||||
ttlMs: number,
|
||||
fn: () => Promise<T>,
|
||||
options: CachedOptions = {},
|
||||
): Promise<T> {
|
||||
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()) {
|
||||
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<T>(
|
||||
key: string,
|
||||
ttlMs: number,
|
||||
staleMs: number,
|
||||
fn: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
const generation = generations.get(key) ?? 0;
|
||||
const compute = (async (): Promise<T> => {
|
||||
// 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;
|
||||
});
|
||||
|
||||
// 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, ttlSec, JSON.stringify(data));
|
||||
await redis.setex(
|
||||
key,
|
||||
redisTtlSeconds(Math.ceil(ttlMs / 1000)),
|
||||
JSON.stringify(data),
|
||||
);
|
||||
} catch {
|
||||
/* non-critical: memory cache still works */
|
||||
}
|
||||
}
|
||||
setMemory(key, { data, expiresAt: Date.now() + ttlMs });
|
||||
reportedFailures.delete(key);
|
||||
commit(key, data, ttlMs, staleMs, generation);
|
||||
}
|
||||
|
||||
return data;
|
||||
})().finally(() => {
|
||||
@@ -99,3 +252,24 @@ export async function cached<T>(
|
||||
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);
|
||||
@@ -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<string>(),
|
||||
}));
|
||||
|
||||
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"]);
|
||||
});
|
||||
+18
-7
@@ -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<T>(
|
||||
cacheKey: string,
|
||||
queryFn: () => Promise<T>,
|
||||
ttl: number = DEFAULT_CACHE_TTL,
|
||||
options?: CachedOptions,
|
||||
): Promise<T> {
|
||||
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<number> {
|
||||
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;
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
+81
-8
@@ -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<string, number>();
|
||||
|
||||
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<void> {
|
||||
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<string, string[]>();
|
||||
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<void> {
|
||||
for (const name of files) {
|
||||
try {
|
||||
await unlink(join(directory, name));
|
||||
} catch {
|
||||
// Already gone.
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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<T>(
|
||||
key: string,
|
||||
ttlSeconds: number,
|
||||
fetch: () => Promise<T>,
|
||||
options?: CachedOptions,
|
||||
): Promise<T> {
|
||||
return cached(key, ttlSeconds * 1000, fetch);
|
||||
return cached(key, ttlSeconds * 1000, fetch, options);
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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<unknown>) =>
|
||||
query(),
|
||||
),
|
||||
redisCache: vi.fn(
|
||||
async (_key: string, _ttl: number, query: () => Promise<unknown>) =>
|
||||
query(),
|
||||
),
|
||||
cacheNews: vi.fn(
|
||||
async (_key: string, _ttl: number, query: () => Promise<unknown>) =>
|
||||
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();
|
||||
});
|
||||
@@ -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<T>(name: string, load: () => Promise<T>): Promise<void> {
|
||||
try {
|
||||
await load();
|
||||
} catch (error) {
|
||||
logger.warn("[cache] warm-up failed", {
|
||||
key: name,
|
||||
error: String(error),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
export async function warmPublicCaches(): Promise<void> {
|
||||
// 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<number> {
|
||||
const [row] = await db
|
||||
.select({ total: count() })
|
||||
.from(User)
|
||||
.where(eq(User.online, "1"));
|
||||
return row?.total ?? 0;
|
||||
}
|
||||
|
||||
async function countUsers(): Promise<number> {
|
||||
const [row] = await db.select({ total: count() }).from(User);
|
||||
return row?.total ?? 0;
|
||||
}
|
||||
|
||||
async function countRooms(): Promise<number> {
|
||||
const [row] = await db.select({ total: count() }).from(Rooms);
|
||||
return row?.total ?? 0;
|
||||
}
|
||||
|
||||
async function countPhotos(): Promise<number> {
|
||||
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));
|
||||
}
|
||||
@@ -6,10 +6,33 @@ const state = vi.hoisted(() => ({
|
||||
values: new Map<string, unknown>(),
|
||||
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<unknown>) => {
|
||||
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();
|
||||
|
||||
@@ -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<string> | null = null;
|
||||
|
||||
async function readRevision(): Promise<string> {
|
||||
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<T>(
|
||||
key: string,
|
||||
ttlMs: number,
|
||||
query: () => Promise<T>,
|
||||
): Promise<T> {
|
||||
if (!redis || redis.status === "end") return query();
|
||||
let revision: string;
|
||||
try {
|
||||
revision = (await redis.get(REVISION_KEY)) ?? "0";
|
||||
} catch {
|
||||
return 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);
|
||||
}
|
||||
return cached(`news:${revision}:${key}`, ttlMs, query);
|
||||
|
||||
async function rotateRevision(): Promise<string> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
if (!redis) return;
|
||||
if (redis.status === "end")
|
||||
throw new Error("News cache connection is closed");
|
||||
await redis.set(REVISION_KEY, randomUUID());
|
||||
await rotateRevision();
|
||||
}
|
||||
Reference in new issue
Block a user