// Distributed lock abstraction with two implementations: // - memory: per-process key set with TTL (only meaningful within one instance) // - redis: SET key token NX PX ttl, released with a compare-and-delete Lua // script so only the holder can release it // // Use acquire/release for long-lived ownership (e.g. a background poller) and // withLock for a one-shot critical section. Selection is based on REDIS_URL. // // There is deliberately no auto-renewal/watchdog: every guarded section (sweep // jobs, template seeding) is a handful of status-conditional bulk UPDATEs that // finish far below the lock TTL, and a rare TTL overrun only risks one // idempotent overlapping run. Keep TTLs generous instead of adding renewal. import { randomUUID } from 'crypto'; import { getRedis, isRedisEnabled } from '../redis.js'; // Thrown when Redis is configured but the lock backend cannot be reached. // This is distinct from contention (acquire resolves null): the caller must // decide whether its critical section is safe to run without mutual exclusion. export class LockUnavailableError extends Error { constructor(message: string) { super(message); this.name = 'LockUnavailableError'; } } export interface WithLockOptions { // What withLock should do when the lock backend is unavailable (not mere // contention): 'skip' (default) returns null as if the lock were held // elsewhere; 'run' executes fn without mutual exclusion. onUnavailable?: 'skip' | 'run'; } export interface Lock { readonly backend: 'memory' | 'redis'; // Returns a token when the lock was acquired, or null when already held. // Throws LockUnavailableError when the backend is configured but erroring. acquire(key: string, ttlMs: number): Promise; release(key: string, token: string): Promise; // Runs fn while holding the lock; returns fn's result, or null if not acquired. withLock(key: string, ttlMs: number, fn: () => Promise, opts?: WithLockOptions): Promise; } // ==================== Memory implementation ==================== export class MemoryLock implements Lock { readonly backend = 'memory' as const; private held = new Map(); async acquire(key: string, ttlMs: number): Promise { const existing = this.held.get(key); const now = Date.now(); if (existing && existing.expiresAt > now) { return null; } const token = randomUUID(); this.held.set(key, { token, expiresAt: now + ttlMs }); return token; } async release(key: string, token: string): Promise { const existing = this.held.get(key); if (existing && existing.token === token) { this.held.delete(key); } } async withLock(key: string, ttlMs: number, fn: () => Promise): Promise { const token = await this.acquire(key, ttlMs); if (!token) return null; try { return await fn(); } finally { await this.release(key, token); } } } // ==================== Redis implementation ==================== const RELEASE_SCRIPT = 'if redis.call("get", KEYS[1]) == ARGV[1] then return redis.call("del", KEYS[1]) else return 0 end'; export class RedisLock implements Lock { readonly backend = 'redis' as const; async acquire(key: string, ttlMs: number): Promise { const redis = getRedis(); if (!redis) { throw new LockUnavailableError('redis client not initialized'); } const token = randomUUID(); try { const result = await redis.set(`lock:${key}`, token, 'PX', ttlMs, 'NX'); return result === 'OK' ? token : null; } catch (err: any) { // Do NOT fabricate a token here: during an outage every replica would // "acquire" every lock and run the guarded sections concurrently. Let the // caller choose between skipping the run and running unlocked. throw new LockUnavailableError(err?.message || String(err)); } } async release(key: string, token: string): Promise { const redis = getRedis(); if (!redis) return; try { await redis.eval(RELEASE_SCRIPT, 1, `lock:${key}`, token); } catch (err: any) { console.error('[lock] redis release error:', err?.message || err); } } async withLock(key: string, ttlMs: number, fn: () => Promise, opts?: WithLockOptions): Promise { let token: string | null; try { token = await this.acquire(key, ttlMs); } catch (err) { if (!(err instanceof LockUnavailableError)) throw err; if (opts?.onUnavailable === 'run') { console.warn(`[lock] backend unavailable, running "${key}" WITHOUT mutual exclusion:`, err.message); return fn(); } console.warn(`[lock] backend unavailable, skipping "${key}" this run:`, err.message); return null; } if (!token) return null; try { return await fn(); } finally { await this.release(key, token); } } } // ==================== Selection ==================== let instance: Lock | null = null; export function getLock(): Lock { if (!instance) { instance = isRedisEnabled() ? new RedisLock() : new MemoryLock(); } return instance; }