54 lines
1.6 KiB
TypeScript
54 lines
1.6 KiB
TypeScript
import type Redis from "ioredis";
|
|
|
|
import { cacheTagsForOutboxEvent } from "@/lib/cache/tags";
|
|
import {
|
|
expectedTopic,
|
|
toInvalidationEvent,
|
|
type InvalidationEvent,
|
|
type OutboxEventLike,
|
|
} from "@/lib/events/invalidation";
|
|
import {
|
|
redisCacheTagStateKey,
|
|
redisEventChannel,
|
|
} from "@/lib/redis/client";
|
|
|
|
export interface OutboxTransport {
|
|
invalidate(tags: string[], timestamp: number): Promise<void>;
|
|
publish(topic: string, event: InvalidationEvent): Promise<void>;
|
|
}
|
|
|
|
export function createRedisOutboxTransport(client: Redis): OutboxTransport {
|
|
return {
|
|
async invalidate(tags, timestamp) {
|
|
if (tags.length === 0) return;
|
|
const pipeline = client.pipeline();
|
|
for (const tag of tags) {
|
|
pipeline.hset(
|
|
redisCacheTagStateKey(),
|
|
tag,
|
|
JSON.stringify({ expired: timestamp }),
|
|
);
|
|
}
|
|
const results = await pipeline.exec();
|
|
const failure = results?.find(([cause]) => cause);
|
|
if (failure?.[0]) throw failure[0];
|
|
},
|
|
async publish(topic, event) {
|
|
await client.publish(redisEventChannel(topic), JSON.stringify(event));
|
|
},
|
|
};
|
|
}
|
|
|
|
export async function processOutboxEvent(
|
|
record: OutboxEventLike,
|
|
transport: OutboxTransport,
|
|
timestamp = Date.now(),
|
|
): Promise<InvalidationEvent> {
|
|
const event = toInvalidationEvent(record);
|
|
if (!event) throw new Error("Outbox event has an invalid public envelope or topic.");
|
|
const tags = cacheTagsForOutboxEvent(record);
|
|
await transport.invalidate(tags, timestamp);
|
|
await transport.publish(expectedTopic(event), event);
|
|
return event;
|
|
}
|