test(events): cover cross-replica Redis delivery
This commit is contained in:
@@ -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<void> {
|
||||||
|
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<string>((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,
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user