diff --git a/tests/redis-replication.integration.test.ts b/tests/redis-replication.integration.test.ts new file mode 100644 index 0000000..d4c5e49 --- /dev/null +++ b/tests/redis-replication.integration.test.ts @@ -0,0 +1,130 @@ +import { randomUUID } from "node:crypto"; + +import Redis from "ioredis"; +import { afterAll, beforeAll, describe, expect, it } from "vitest"; + +import { 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("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, + }); + }); +});