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

90 lines
2.5 KiB
TypeScript

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