90 lines
2.5 KiB
TypeScript
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")));
|
|
}
|