import type { AddOptions, Job, JobHandler } from "./index.ts"; /** * Persistence contract for the durable queue. * * Distributed drivers should implement `claim()` atomically and exclude leased * jobs from `due()` until their lease expires. Calling `put()` must replace the * stored record and release any previous lease for that job. */ export interface QueueStore { put(job: Job): Promise; get(id: string): Promise; remove(id: string): Promise; due(now: number, limit: number): Promise; list(name?: string): Promise; claim?(id: string, worker: string, leaseUntil: number): Promise; } export function memoryQueueStore(): QueueStore { const jobs = new Map(); return { async put(job) { jobs.set(job.id, structuredClone(job)); }, async get(id) { const job = jobs.get(id); return job ? structuredClone(job) : null; }, async remove(id) { jobs.delete(id); }, async due(now, limit) { return [...jobs.values()] .filter((job) => job.runAt <= now) .sort( (left, right) => right.priority - left.priority || left.runAt - right.runAt || left.createdAt - right.createdAt, ) .slice(0, limit) .map((job) => structuredClone(job)); }, async list(name) { return [...jobs.values()] .filter((job) => !name || job.name === name) .map((job) => structuredClone(job)); }, }; } export interface DurableQueueOptions { store?: QueueStore; workerId?: string; maxAttempts?: number; concurrency?: number; leaseMs?: number; backoff?: (attempt: number) => number; now?: () => number; onDeadLetter?: (job: Job, error: unknown) => void | Promise; } export interface DurableQueue { add(name: string, data: T, options?: AddOptions): Promise>; process(name: string, handler: JobHandler): void; drain(): Promise; failed(): Job[]; retry(id: string): Promise; } function positiveInteger(value: number, label: string): number { if (!Number.isInteger(value) || value < 1) { throw new RangeError(`${label} must be a positive integer`); } return value; } function nonNegativeNumber(value: number, label: string): number { if (!Number.isFinite(value) || value < 0) { throw new RangeError(`${label} must be a non-negative number`); } return value; } export function createDurableQueue(options: DurableQueueOptions = {}): DurableQueue { const store = options.store ?? memoryQueueStore(); const handlers = new Map(); const deadLetters = new Map(); const now = options.now ?? Date.now; const workerId = options.workerId?.trim() || `worker-${crypto.randomUUID()}`; const defaultMaxAttempts = positiveInteger(options.maxAttempts ?? 3, "queue maxAttempts"); const concurrency = positiveInteger(options.concurrency ?? 10, "queue concurrency"); const leaseMs = positiveInteger(options.leaseMs ?? 30_000, "queue leaseMs"); let sequence = 0; let draining = false; async function runJob(job: Job): Promise { const handler = handlers.get(job.name); if (!handler) return false; if (store.claim && !(await store.claim(job.id, workerId, now() + leaseMs))) { return false; } job.attempts += 1; try { await handler(job); if (job.repeat && job.repeat > 0) { job.attempts = 0; job.runAt = now() + job.repeat; await store.put(job); } else { await store.remove(job.id); } } catch (error) { if (job.attempts < job.maxAttempts) { const delay = nonNegativeNumber( options.backoff?.(job.attempts) ?? Math.min(60_000, 1_000 * 2 ** (job.attempts - 1)), "queue retry backoff", ); job.runAt = now() + delay; await store.put(job); } else { await store.remove(job.id); deadLetters.set(job.id, structuredClone(job)); await options.onDeadLetter?.(structuredClone(job), error); } } return true; } return { async add(name: string, data: T, add: AddOptions = {}) { if (!name.trim()) throw new TypeError("queue job name cannot be empty"); const maxAttempts = positiveInteger(add.maxAttempts ?? defaultMaxAttempts, "job maxAttempts"); const delayMs = nonNegativeNumber(add.delayMs ?? 0, "job delayMs"); const priority = add.priority ?? 0; if (!Number.isFinite(priority)) { throw new RangeError("job priority must be a finite number"); } if (add.repeat !== undefined) { positiveInteger(add.repeat, "job repeat"); } if (add.idempotencyKey) { const existing = (await store.list(name)).find( (job) => job.idempotencyKey === add.idempotencyKey, ); if (existing) return existing as Job; } const createdAt = now(); const job: Job = { id: `job-${createdAt.toString(36)}-${(++sequence).toString(36)}`, name, data, attempts: 0, maxAttempts, runAt: createdAt + delayMs, repeat: add.repeat, priority, idempotencyKey: add.idempotencyKey, createdAt, }; await store.put(job); return structuredClone(job); }, process(name, handler) { if (!name.trim()) throw new TypeError("queue worker name cannot be empty"); handlers.set(name, handler as JobHandler); }, async drain() { if (draining) return 0; draining = true; try { const due = await store.due(now(), concurrency); const results = await Promise.all(due.map(runJob)); return results.filter(Boolean).length; } finally { draining = false; } }, failed() { return [...deadLetters.values()].map((job) => structuredClone(job)); }, async retry(id) { const job = deadLetters.get(id); if (!job) return false; deadLetters.delete(id); job.attempts = 0; job.runAt = now(); await store.put(job); return true; }, }; }