diff --git a/lib/events/invalidation.ts b/lib/events/invalidation.ts new file mode 100644 index 0000000..a1c439b --- /dev/null +++ b/lib/events/invalidation.ts @@ -0,0 +1,83 @@ +export const INVALIDATION_TYPES = [ + "directory.updated", + "page.updated", + "data-source.updated", + "admin.updated", +] as const; + +export type InvalidationType = (typeof INVALIDATION_TYPES)[number]; + +export interface InvalidationEvent { + type: InvalidationType; + id: string; + version: number; +} + +export interface OutboxEventLike { + aggregateId: string; + eventType: string; + payload: Record; + topic: string; +} + +function isInvalidationType(value: string): value is InvalidationType { + return INVALIDATION_TYPES.some((type) => type === value); +} + +export function expectedTopic(event: InvalidationEvent): string { + switch (event.type) { + case "directory.updated": + return "directory"; + case "page.updated": + return `page:${event.id}`; + case "data-source.updated": + return `data-source:${event.id}`; + case "admin.updated": + return "admin"; + } +} + +export function toInvalidationEvent( + record: OutboxEventLike, +): InvalidationEvent | null { + if (!isInvalidationType(record.eventType)) return null; + const id = record.payload.id; + const version = record.payload.version; + if ( + typeof id !== "string" || + id !== record.aggregateId || + typeof version !== "number" || + !Number.isSafeInteger(version) || + version < 0 + ) { + return null; + } + const event: InvalidationEvent = { + type: record.eventType, + id, + version, + }; + return expectedTopic(event) === record.topic ? event : null; +} + +export function parseInvalidationEvent(value: string): InvalidationEvent | null { + try { + const parsed = JSON.parse(value) as Record; + const type = parsed.type; + const id = parsed.id; + const version = parsed.version; + if ( + typeof type !== "string" || + !isInvalidationType(type) || + typeof id !== "string" || + typeof version !== "number" || + !Number.isSafeInteger(version) || + version < 0 + ) { + return null; + } + return { type, id, version }; + } catch { + return null; + } +} diff --git a/lib/outbox/processor.test.ts b/lib/outbox/processor.test.ts new file mode 100644 index 0000000..66b772b --- /dev/null +++ b/lib/outbox/processor.test.ts @@ -0,0 +1,62 @@ +import { describe, expect, it, vi } from "vitest"; + +import { processOutboxEvent, type OutboxTransport } from "./processor"; + +function record(overrides: Partial[0]> = {}) { + return { + topic: "page:72df08ab-50dd-4cbd-9a69-70949d34cf9f", + aggregateId: "72df08ab-50dd-4cbd-9a69-70949d34cf9f", + eventType: "page.updated", + payload: { + id: "72df08ab-50dd-4cbd-9a69-70949d34cf9f", + version: 8, + privateContent: "must not leave PostgreSQL", + }, + ...overrides, + }; +} + +describe("transactional outbox processor", () => { + it("invalidates before publishing a privacy-safe event", async () => { + const calls: string[] = []; + const transport: OutboxTransport = { + invalidate: vi.fn(async () => { + calls.push("invalidate"); + }), + publish: vi.fn(async () => { + calls.push("publish"); + }), + }; + + const event = await processOutboxEvent(record(), transport, 1234); + + expect(calls).toEqual(["invalidate", "publish"]); + expect(transport.invalidate).toHaveBeenCalledWith( + ["buzz:page:72df08ab-50dd-4cbd-9a69-70949d34cf9f"], + 1234, + ); + expect(transport.publish).toHaveBeenCalledWith( + "page:72df08ab-50dd-4cbd-9a69-70949d34cf9f", + { + type: "page.updated", + id: "72df08ab-50dd-4cbd-9a69-70949d34cf9f", + version: 8, + }, + ); + expect(event).not.toHaveProperty("privateContent"); + }); + + it("rejects a mismatched topic or malformed version", async () => { + const transport: OutboxTransport = { + invalidate: vi.fn(), + publish: vi.fn(), + }; + await expect( + processOutboxEvent(record({ topic: "admin" }), transport), + ).rejects.toThrow("invalid public envelope"); + await expect( + processOutboxEvent(record({ payload: { id: record().aggregateId, version: "8" } }), transport), + ).rejects.toThrow("invalid public envelope"); + expect(transport.invalidate).not.toHaveBeenCalled(); + }); +}); diff --git a/lib/outbox/processor.ts b/lib/outbox/processor.ts new file mode 100644 index 0000000..8a3afbc --- /dev/null +++ b/lib/outbox/processor.ts @@ -0,0 +1,53 @@ +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; + publish(topic: string, event: InvalidationEvent): Promise; +} + +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 { + 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; +} diff --git a/lib/outbox/repository.test.ts b/lib/outbox/repository.test.ts new file mode 100644 index 0000000..8bc92bf --- /dev/null +++ b/lib/outbox/repository.test.ts @@ -0,0 +1,11 @@ +import { describe, expect, it } from "vitest"; + +import { retryDelayMs } from "./repository"; + +describe("outbox retry policy", () => { + it("uses capped exponential backoff", () => { + expect(retryDelayMs(1)).toBe(500); + expect(retryDelayMs(4)).toBe(4_000); + expect(retryDelayMs(30)).toBe(60_000); + }); +}); diff --git a/lib/outbox/repository.ts b/lib/outbox/repository.ts new file mode 100644 index 0000000..eb6acc8 --- /dev/null +++ b/lib/outbox/repository.ts @@ -0,0 +1,89 @@ +import { and, asc, eq, inArray, lte, sql } from "drizzle-orm"; + +import { getDb } from "@/db/client"; +import { outboxEvents } from "@/db/schema"; + +export type ClaimedOutboxEvent = Pick< + typeof outboxEvents.$inferSelect, + "id" | "topic" | "aggregateId" | "eventType" | "payload" | "attempts" +>; + +const CLAIMABLE_STATUSES = ["pending", "processing"] as const; + +export async function claimOutboxEvents( + batchSize = 20, + leaseMs = 30_000, + now = new Date(), +): Promise { + const size = Math.max(1, Math.min(100, Math.floor(batchSize))); + const leaseUntil = new Date(now.getTime() + Math.max(5_000, leaseMs)); + + return getDb().transaction(async (tx) => { + const candidates = await tx + .select({ id: outboxEvents.id }) + .from(outboxEvents) + .where( + and( + inArray(outboxEvents.status, [...CLAIMABLE_STATUSES]), + lte(outboxEvents.availableAt, now), + ), + ) + .orderBy(asc(outboxEvents.id)) + .limit(size) + .for("update", { skipLocked: true }); + + const ids = candidates.map(({ id }) => id); + if (ids.length === 0) return []; + return tx + .update(outboxEvents) + .set({ + status: "processing", + attempts: sql`${outboxEvents.attempts} + 1`, + availableAt: leaseUntil, + }) + .where(inArray(outboxEvents.id, ids)) + .returning({ + id: outboxEvents.id, + topic: outboxEvents.topic, + aggregateId: outboxEvents.aggregateId, + eventType: outboxEvents.eventType, + payload: outboxEvents.payload, + attempts: outboxEvents.attempts, + }); + }); +} + +export async function markOutboxProcessed( + id: number, + now = new Date(), +): Promise { + await getDb() + .update(outboxEvents) + .set({ + status: "processed", + processedAt: now, + lastError: null, + }) + .where(and(eq(outboxEvents.id, id), eq(outboxEvents.status, "processing"))); +} + +export function retryDelayMs(attempts: number): number { + return Math.min(60_000, 500 * 2 ** Math.max(0, Math.min(attempts - 1, 16))); +} + +export async function releaseOutboxEvent( + id: number, + attempts: number, + cause: unknown, + now = new Date(), +): Promise { + const message = cause instanceof Error ? cause.message : String(cause); + await getDb() + .update(outboxEvents) + .set({ + status: "pending", + availableAt: new Date(now.getTime() + retryDelayMs(attempts)), + lastError: message.slice(0, 2_000), + }) + .where(and(eq(outboxEvents.id, id), eq(outboxEvents.status, "processing"))); +} diff --git a/lib/redis/client.ts b/lib/redis/client.ts new file mode 100644 index 0000000..ad490c7 --- /dev/null +++ b/lib/redis/client.ts @@ -0,0 +1,90 @@ +import Redis from "ioredis"; + +let sharedClient: Redis | undefined; +let sharedConnection: Promise | undefined; + +function redisUrl(): string { + const value = process.env.REDIS_URL; + if (!value) throw new Error("REDIS_URL is required for Redis-backed features."); + return value; +} + +export function createRedisClient(connectionName: string): Redis { + const client = new Redis(redisUrl(), { + connectionName, + enableOfflineQueue: false, + lazyConnect: true, + maxRetriesPerRequest: 2, + }); + client.on("error", () => undefined); + return client; +} + +export async function connectRedisClient(client: Redis): Promise { + if (client.status === "ready") return client; + if (client.status === "wait") { + await client.connect(); + return client; + } + + await new Promise((resolve, reject) => { + const ready = () => { + cleanup(); + resolve(); + }; + const error = (cause: Error) => { + cleanup(); + reject(cause); + }; + const cleanup = () => { + client.off("ready", ready); + client.off("error", error); + }; + client.once("ready", ready); + client.once("error", error); + }); + return client; +} + +export async function getRedisClient(): Promise { + if (!sharedClient) { + sharedClient = createRedisClient("buzz-sheet-shared"); + } + if (!sharedConnection) { + sharedConnection = connectRedisClient(sharedClient).catch((cause) => { + sharedClient?.disconnect(); + sharedClient = undefined; + sharedConnection = undefined; + throw cause; + }); + } + return sharedConnection; +} + +export async function closeRedisClient(): Promise { + const client = sharedClient; + sharedClient = undefined; + sharedConnection = undefined; + if (!client) return; + if (client.status === "wait" || client.status === "end") { + client.disconnect(); + return; + } + await client.quit().catch(() => client.disconnect()); +} + +export function redisCachePrefix(): string { + return process.env.REDIS_CACHE_PREFIX || "buzz:next-cache"; +} + +export function redisCacheTagStateKey(): string { + return `${redisCachePrefix()}:tag-timestamps`; +} + +export function redisEventPrefix(): string { + return process.env.REDIS_EVENT_PREFIX || "buzz:events"; +} + +export function redisEventChannel(topic: string): string { + return `${redisEventPrefix()}:${topic}`; +} diff --git a/package.json b/package.json index fde0f44..19d158d 100644 --- a/package.json +++ b/package.json @@ -14,7 +14,8 @@ "db:generate": "drizzle-kit generate", "db:migrate": "bun scripts/migrate.ts", "db:studio": "drizzle-kit studio", - "db:seed": "bun scripts/seed-templates.ts" + "db:seed": "bun scripts/seed-templates.ts", + "worker:outbox": "bun scripts/outbox-worker.ts" }, "dependencies": { "@base-ui/react": "^1.7.0", diff --git a/scripts/outbox-worker.ts b/scripts/outbox-worker.ts new file mode 100644 index 0000000..f8798a8 --- /dev/null +++ b/scripts/outbox-worker.ts @@ -0,0 +1,65 @@ +import { closeDb } from "@/db/client"; +import { + claimOutboxEvents, + markOutboxProcessed, + releaseOutboxEvent, +} from "@/lib/outbox/repository"; +import { + createRedisOutboxTransport, + processOutboxEvent, +} from "@/lib/outbox/processor"; +import { closeRedisClient, getRedisClient } from "@/lib/redis/client"; + +let stopping = false; +const stop = () => { + stopping = true; +}; +process.once("SIGINT", stop); +process.once("SIGTERM", stop); + +function wait(milliseconds: number): Promise { + return new Promise((resolve) => setTimeout(resolve, milliseconds)); +} + +function positiveInteger( + name: string, + fallback: number, + maximum: number, +): number { + const parsed = Number(process.env[name]); + return Number.isSafeInteger(parsed) && parsed > 0 + ? Math.min(parsed, maximum) + : fallback; +} + +async function run(): Promise { + const redis = await getRedisClient(); + const transport = createRedisOutboxTransport(redis); + const batchSize = positiveInteger("OUTBOX_BATCH_SIZE", 20, 100); + const pollMs = positiveInteger("OUTBOX_POLL_MS", 1_000, 60_000); + const leaseMs = positiveInteger("OUTBOX_LEASE_MS", 30_000, 300_000); + + while (!stopping) { + const events = await claimOutboxEvents(batchSize, leaseMs); + if (events.length === 0) { + await wait(pollMs); + continue; + } + await Promise.all( + events.map(async (event) => { + try { + await processOutboxEvent(event, transport); + await markOutboxProcessed(event.id); + } catch (cause) { + await releaseOutboxEvent(event.id, event.attempts, cause); + } + }), + ); + } +} + +try { + await run(); +} finally { + await Promise.all([closeDb(), closeRedisClient()]); +}