41 lines
1.7 KiB
TypeScript
41 lines
1.7 KiB
TypeScript
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<number> {
|
|
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;
|
|
}
|