import { and, asc, eq, isNull, lte, notExists } from "drizzle-orm"; import { getDb } from "@/db/client"; import { commissionCheckouts, commissionTickets } from "@/db/schema"; import { CHECKOUT_DURATION_MS } from "@/lib/commission/checkout-expiration"; import { commissionSlipLockKey } from "@/lib/commission/slip-lock"; import { getRedisClient } from "@/lib/redis/client"; const BATCH_SIZE = 100; export async function removeExpiredUnpaidCheckouts(): Promise { const db = getDb(); const redis = await getRedisClient(); const cutoff = new Date(Date.now() - CHECKOUT_DURATION_MS); const candidates = await db.select({ id: commissionCheckouts.id }) .from(commissionCheckouts) .leftJoin(commissionTickets, eq(commissionTickets.checkoutId, commissionCheckouts.id)) .where(and(lte(commissionCheckouts.createdAt, cutoff), isNull(commissionTickets.id))) .orderBy(asc(commissionCheckouts.createdAt)) .limit(BATCH_SIZE); let removed = 0; for (const { id } of candidates) { const lockKey = commissionSlipLockKey(id); const lockId = crypto.randomUUID(); if (await redis.set(lockKey, lockId, "EX", 60, "NX") !== "OK") continue; try { const deleted = await db.delete(commissionCheckouts).where(and( eq(commissionCheckouts.id, id), lte(commissionCheckouts.createdAt, cutoff), notExists(db.select({ id: commissionTickets.id }).from(commissionTickets) .where(eq(commissionTickets.checkoutId, id))), )).returning({ id: commissionCheckouts.id }); removed += deleted.length; } finally { await redis.eval("if redis.call('GET', KEYS[1]) == ARGV[1] then return redis.call('DEL', KEYS[1]) end return 0", 1, lockKey, lockId).catch(() => undefined); } } return removed; }