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