import { createExecutionContext, type ExecutionContext } from "@wrnexus/core"; 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; findByIdempotencyKey?(name: string, key: string): Promise; size?(): Promise; claim?(id: string, worker: string, leaseUntil: number, now?: number): Promise; release?(id: string, worker: string): Promise; archive?(record: QueueJobRecord): Promise; history?(id?: string): Promise; } export type QueueJobState = "completed" | "failed" | "cancelled"; export interface QueueJobRecord { job: Job; state: QueueJobState; finishedAt: number; error?: string; } export type QueueEventType = | "job.added" | "job.started" | "job.completed" | "job.retrying" | "job.failed" | "job.cancelled" | "worker.started" | "worker.stopped" | "worker.error"; export interface QueueEvent { type: QueueEventType; at: number; workerId: string; job?: Job; error?: unknown; durationMs?: number; } export interface QueueHealth { running: boolean; accepting: boolean; draining: boolean; active: number; pending: number; failed: number; workerId: string; lastPollAt?: number; lastError?: unknown; } export function memoryQueueStore(): QueueStore { const jobs = new Map(); const records = 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)); }, async findByIdempotencyKey(name, key) { const job = [...jobs.values()].find( (candidate) => candidate.name === name && candidate.idempotencyKey === key, ); return job ? structuredClone(job) : null; }, async size() { return jobs.size; }, async archive(record) { records.set(record.job.id, structuredClone(record)); }, async history(id) { return [...records.values()] .filter((record) => !id || record.job.id === id) .map((record) => structuredClone(record)); }, }; } export interface DurableQueueOptions { store?: QueueStore; workerId?: string; maxAttempts?: number; concurrency?: number; capacity?: number; leaseMs?: number; pollMs?: number; backoff?: (attempt: number) => number; now?: () => number; onDeadLetter?: (job: Job, error: unknown) => void | Promise; context?: (job: Job, signal: AbortSignal) => ExecutionContext; onEvent?: (event: QueueEvent) => void | Promise; /** Maintenance hook run before each poll, for application-specific recovery. */ beforeDrain?: (queue: DurableQueue) => void | Promise; } export interface DurableQueue { add(name: string, data: T, options?: AddOptions): Promise>; process(name: string, handler: JobHandler): void; drain(): Promise; start(): void; stop(): void; isRunning(): boolean; health(): Promise; get(id: string): Promise; status(id: string): Promise<"queued" | QueueJobState | "missing">; history(id?: string): Promise; cancel(id: string): Promise; shutdown(options?: { force?: boolean }): Promise; list(name?: string): Promise; failed(): Job[]; retry(id: string): Promise; addBatch( name: string, entries: readonly { data: T; options?: AddOptions }[], ): 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 capacity = positiveInteger(options.capacity ?? 10_000, "queue capacity"); const leaseMs = positiveInteger(options.leaseMs ?? 30_000, "queue leaseMs"); const pollMs = positiveInteger(options.pollMs ?? 250, "queue pollMs"); let sequence = 0; let draining = false; let accepting = true; let timer: ReturnType | null = null; let lastPollAt: number | undefined; let lastError: unknown; const active = new Map }>(); const emit = async (event: Omit) => { await options.onEvent?.({ ...event, at: now(), workerId }); }; function beginJob(job: Job): Promise { const controller = new AbortController(); const promise = runJob(job, controller).finally(() => active.delete(job.id)); active.set(job.id, { controller, promise }); return promise; } async function runJob(job: Job, controller: AbortController): 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; const startedAt = now(); await emit({ type: "job.started", job: structuredClone(job) }); try { await handler(job, { signal: controller.signal, execution: options.context?.(job, controller.signal) ?? createExecutionContext({ kind: job.repeat ? "cron" : "queue", signal: controller.signal, metadata: { jobId: job.id, jobName: job.name, attempt: job.attempts }, }), }); if (job.repeat && job.repeat > 0) { job.attempts = 0; job.runAt = now() + job.repeat; await store.put(job); } else { await store.archive?.({ job: structuredClone(job), state: "completed", finishedAt: now() }); await store.remove(job.id); } await emit({ type: "job.completed", job: structuredClone(job), durationMs: now() - startedAt, }); } catch (error) { if (controller.signal.aborted) { await store.remove(job.id); return true; } 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); await emit({ type: "job.retrying", job: structuredClone(job), error }); } else { await store.remove(job.id); deadLetters.set(job.id, structuredClone(job)); await store.archive?.({ job: structuredClone(job), state: "failed", finishedAt: now(), error: error instanceof Error ? error.message : String(error), }); await options.onDeadLetter?.(structuredClone(job), error); await emit({ type: "job.failed", job: structuredClone(job), error }); } } return true; } const durableQueue: DurableQueue = { async add(name: string, data: T, add: AddOptions = {}) { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); 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 = store.findByIdempotencyKey ? await store.findByIdempotencyKey(name, add.idempotencyKey) : (await store.list(name)).find((job) => job.idempotencyKey === add.idempotencyKey); if (existing) return existing as Job; } if ((store.size ? await store.size() : (await store.list()).length) >= capacity) { throw new Error(`WRN-QUEUE-CAPACITY: queue capacity of ${capacity} reached`); } 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); await emit({ type: "job.added", job: structuredClone(job) }); return structuredClone(job); }, async addBatch(name: string, entries: readonly { data: T; options?: AddOptions }[]) { const output: Job[] = []; for (const entry of entries) output.push(await this.add(name, entry.data, entry.options)); return output; }, 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; lastPollAt = now(); try { await options.beforeDrain?.(durableQueue); const due = await store.due(now(), concurrency); const results = await Promise.all(due.map(beginJob)); return results.filter(Boolean).length; } finally { draining = false; } }, failed() { return [...deadLetters.values()].map((job) => structuredClone(job)); }, list(name) { return store.list(name); }, get(id) { return store.get(id); }, async status(id) { if (await store.get(id)) return "queued"; const record = (await store.history?.(id))?.at(-1); if (record) return record.state; if (deadLetters.has(id)) return "failed"; return "missing"; }, async history(id) { return (await store.history?.(id)) ?? []; }, start() { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); if (timer) return; timer = setInterval(() => { void this.drain().catch(async (error: unknown) => { lastError = error; await emit({ type: "worker.error", error }); }); }, pollMs); timer.unref?.(); void emit({ type: "worker.started" }); }, stop() { if (timer) clearInterval(timer); timer = null; void emit({ type: "worker.stopped" }); }, isRunning() { return timer !== null; }, async health() { return { running: timer !== null, accepting, draining, active: active.size, pending: store.size ? await store.size() : (await store.list()).length, failed: deadLetters.size, workerId, lastPollAt, lastError, }; }, async cancel(id) { const running = active.get(id); if (running) { running.controller.abort(); const job = await store.get(id); if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() }); await emit({ type: "job.cancelled" }); return true; } if (!(await store.get(id))) return false; const job = await store.get(id); if (job) await store.archive?.({ job, state: "cancelled", finishedAt: now() }); await store.remove(id); await emit({ type: "job.cancelled" }); return true; }, async shutdown(shutdownOptions = {}) { accepting = false; if (timer) clearInterval(timer); timer = null; if (shutdownOptions.force) { for (const { controller } of active.values()) controller.abort(); } await Promise.allSettled([...active.values()].map(({ promise }) => promise)); }, async retry(id) { let job = deadLetters.get(id); if (!job) { const record = (await store.history?.(id))?.find( (candidate) => candidate.state === "failed", ); job = record?.job; } if (!job) return false; deadLetters.delete(id); job.attempts = 0; job.runAt = now(); await store.put(job); return true; }, }; return durableQueue; }