58 lines
1.4 KiB
TypeScript
58 lines
1.4 KiB
TypeScript
import { randomUUID } from "node:crypto";
|
|
|
|
import { getRedisClient, redisEventPrefix } from "@/lib/redis/client";
|
|
|
|
const LOCK_TTL_MS = 120_000;
|
|
const RENEW_INTERVAL_MS = 30_000;
|
|
|
|
const renewScript = `
|
|
if redis.call("get", KEYS[1]) == ARGV[1] then
|
|
return redis.call("pexpire", KEYS[1], ARGV[2])
|
|
end
|
|
return 0
|
|
`;
|
|
|
|
const releaseScript = `
|
|
if redis.call("get", KEYS[1]) == ARGV[1] then
|
|
return redis.call("del", KEYS[1])
|
|
end
|
|
return 0
|
|
`;
|
|
|
|
class CatalogSyncBusyError extends Error {
|
|
constructor() {
|
|
super("กำลัง Sync Catalog อยู่");
|
|
this.name = "CatalogSyncBusyError";
|
|
}
|
|
}
|
|
|
|
export async function withCatalogSyncLock<T>(
|
|
task: () => Promise<T>,
|
|
): Promise<T> {
|
|
const redis = await getRedisClient();
|
|
const key = `${redisEventPrefix()}:catalog-sync-lock`;
|
|
const token = randomUUID();
|
|
const acquired = await redis.set(key, token, "PX", LOCK_TTL_MS, "NX");
|
|
if (acquired !== "OK") throw new CatalogSyncBusyError();
|
|
|
|
let renewing = false;
|
|
const timer = setInterval(() => {
|
|
if (renewing) return;
|
|
renewing = true;
|
|
void redis
|
|
.eval(renewScript, 1, key, token, String(LOCK_TTL_MS))
|
|
.catch(() => undefined)
|
|
.finally(() => {
|
|
renewing = false;
|
|
});
|
|
}, RENEW_INTERVAL_MS);
|
|
timer.unref();
|
|
|
|
try {
|
|
return await task();
|
|
} finally {
|
|
clearInterval(timer);
|
|
await redis.eval(releaseScript, 1, key, token).catch(() => undefined);
|
|
}
|
|
}
|