feat(events): process transactional invalidations
This commit is contained in:
@@ -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<string, unknown>;
|
||||
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<string, unknown>;
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
import { describe, expect, it, vi } from "vitest";
|
||||
|
||||
import { processOutboxEvent, type OutboxTransport } from "./processor";
|
||||
|
||||
function record(overrides: Partial<Parameters<typeof processOutboxEvent>[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();
|
||||
});
|
||||
});
|
||||
@@ -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<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;
|
||||
}
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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<ClaimedOutboxEvent[]> {
|
||||
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<number>`${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<void> {
|
||||
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<void> {
|
||||
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")));
|
||||
}
|
||||
@@ -0,0 +1,90 @@
|
||||
import Redis from "ioredis";
|
||||
|
||||
let sharedClient: Redis | undefined;
|
||||
let sharedConnection: Promise<Redis> | 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<Redis> {
|
||||
if (client.status === "ready") return client;
|
||||
if (client.status === "wait") {
|
||||
await client.connect();
|
||||
return client;
|
||||
}
|
||||
|
||||
await new Promise<void>((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<Redis> {
|
||||
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<void> {
|
||||
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}`;
|
||||
}
|
||||
+2
-1
@@ -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",
|
||||
|
||||
@@ -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<void> {
|
||||
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<void> {
|
||||
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()]);
|
||||
}
|
||||
Reference in New Issue
Block a user