Files
buzz-sheet/lib/outbox/processor.ts
T

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;
}