import { randomUUID } from "node:crypto"; import Redis from "ioredis"; import { afterAll, beforeAll, describe, expect, it } from "vitest"; import { CATALOG_CACHE_TAG, GLOSSARY_CACHE_TAG, pageCacheTag } from "@/lib/cache/tags"; import { parseInvalidationEvent } from "@/lib/events/invalidation"; import { createRedisOutboxTransport, processOutboxEvent, } from "@/lib/outbox/processor"; import { redisCacheTagStateKey, redisEventChannel, } from "@/lib/redis/client"; const integrationUrl = process.env.REDIS_INTEGRATION_URL; const describeWithRedis = integrationUrl ? describe : describe.skip; function integrationClient(name: string): Redis { const client = new Redis(integrationUrl!, { connectionName: name, enableOfflineQueue: false, lazyConnect: true, maxRetriesPerRequest: 1, }); client.on("error", () => undefined); return client; } async function closeClient(client: Redis | undefined): Promise { if (!client) return; if (client.status === "wait" || client.status === "end") { client.disconnect(); return; } await client.quit().catch(() => client.disconnect()); } describeWithRedis("cross-replica Redis invalidation delivery", () => { const runId = randomUUID(); const originalCachePrefix = process.env.REDIS_CACHE_PREFIX; const originalEventPrefix = process.env.REDIS_EVENT_PREFIX; const pageId = `integration-${runId}`; const topic = `page:${pageId}`; const cacheTag = pageCacheTag(pageId); const invalidatedAt = 1_725_000_000_000; let replicaA: Redis; let replicaBEvents: Redis; let replicaBCache: Redis; beforeAll(async () => { process.env.REDIS_CACHE_PREFIX = `buzz:test:${runId}:cache`; process.env.REDIS_EVENT_PREFIX = `buzz:test:${runId}:events`; replicaA = integrationClient(`buzz-sheet-test-a-${runId}`); replicaBEvents = integrationClient(`buzz-sheet-test-b-events-${runId}`); replicaBCache = integrationClient(`buzz-sheet-test-b-cache-${runId}`); await Promise.all([ replicaA.connect(), replicaBEvents.connect(), replicaBCache.connect(), ]); }); afterAll(async () => { if (replicaA?.status === "ready") { await replicaA.del(redisCacheTagStateKey()); } await Promise.all([ closeClient(replicaA), closeClient(replicaBEvents), closeClient(replicaBCache), ]); if (originalCachePrefix === undefined) delete process.env.REDIS_CACHE_PREFIX; else process.env.REDIS_CACHE_PREFIX = originalCachePrefix; if (originalEventPrefix === undefined) delete process.env.REDIS_EVENT_PREFIX; else process.env.REDIS_EVENT_PREFIX = originalEventPrefix; }); it.each([ ["catalog.updated", CATALOG_CACHE_TAG], ["glossary.updated", GLOSSARY_CACHE_TAG], ])("expires %s data across replicas", async (eventType, tag) => { await processOutboxEvent({ topic: "admin", aggregateId: pageId, eventType, payload: { id: pageId, version: 43 }, }, createRedisOutboxTransport(replicaA), invalidatedAt); expect(JSON.parse((await replicaBCache.hget(redisCacheTagStateKey(), tag))!)) .toEqual({ expired: invalidatedAt }); }); it("shares tag state and publishes only the typed event envelope", async () => { const channel = redisEventChannel(topic); await replicaBEvents.subscribe(channel); const received = new Promise((resolve, reject) => { const timeout = setTimeout( () => reject(new Error("Timed out waiting for the replica event.")), 5_000, ); replicaBEvents.on("message", (receivedChannel, value) => { if (receivedChannel !== channel) return; clearTimeout(timeout); resolve(value); }); }); const processed = await processOutboxEvent( { aggregateId: pageId, eventType: "page.updated", payload: { id: pageId, version: 42, privateContent: "must remain in PostgreSQL", }, topic, }, createRedisOutboxTransport(replicaA), invalidatedAt, ); const rawEvent = await received; expect(parseInvalidationEvent(rawEvent)).toEqual(processed); expect(JSON.parse(rawEvent)).toEqual({ type: "page.updated", id: pageId, version: 42, }); const tagState = await replicaBCache.hget( redisCacheTagStateKey(), cacheTag, ); expect(tagState && JSON.parse(tagState)).toEqual({ expired: invalidatedAt, }); }); });