/** * @wrnexus/queue — a background job queue with delays, retries + backoff, and * concurrent workers. The default store is in-process; a pluggable driver lets * you back it with Redis/SQL for durability across restarts. * * const queue = createQueue(); * queue.process("email", async (job) => { await send(job.data); }); * await queue.add("email", { to: "a@b.com" }, { delayMs: 5000, maxAttempts: 3 }); * queue.start(); // begin polling; queue.stop() to halt * * Tests can drive it deterministically with `await queue.drain(now)`. */ export interface Job { id: string; name: string; data: T; attempts: number; maxAttempts: number; runAt: number; /** If set, re-enqueue this job this many ms after each successful run. */ repeat?: number; priority: number; idempotencyKey?: string; createdAt: number; } export interface JobContext { /** Aborted when an active job is cancelled or the queue is force-stopped. */ signal: AbortSignal; /** The same trusted context shape used by HTTP, actions, realtime and webhooks. */ execution: ExecutionContext; } export type JobHandler = (job: Job, context: JobContext) => void | Promise; export interface AddOptions { /** Delay before the job becomes runnable (ms). */ delayMs?: number; /** Max attempts before it's dead-lettered. Default from queue options. */ maxAttempts?: number; /** Re-enqueue this job this many ms after each successful run (recurring). */ repeat?: number; /** Higher-priority jobs run first when multiple jobs are due. */ priority?: number; /** Prevent duplicate queued work with the same stable key. */ idempotencyKey?: string; } export interface QueueOptions { /** Default max attempts per job. Default 3. */ maxAttempts?: number; /** Base retry backoff (ms); doubles per attempt. Default 1000. */ backoffMs?: number; /** Poll interval when started (ms). Default 250. */ pollMs?: number; /** Called when a job exhausts its attempts. */ onFailed?: (job: Job, error: unknown) => void; /** Maximum jobs executed in one drain. Default: unlimited. */ concurrency?: number; /** Maximum queued + active jobs. Adds reject once this limit is reached. */ capacity?: number; /** Clock injection (tests). Default Date.now. */ now?: () => number; context?: (job: Job, signal: AbortSignal) => ExecutionContext; } export interface Queue { add(name: string, data: T, options?: AddOptions): Promise>; process(name: string, handler: JobHandler): void; /** Run every job whose runAt ≤ now, once. Returns how many ran. */ drain(now?: number): Promise; start(): void; stop(): void; /** Stop accepting work and wait for active handlers (or abort them). */ shutdown(options?: { force?: boolean }): Promise; size(): number; get(id: string): Job | undefined; list(name?: string): Job[]; cancel(id: string): boolean; failed(): Job[]; retry(id: string): Promise; } export interface JobDefinition { name: string; options?: Omit; run: JobHandler; } export function defineJob(definition: JobDefinition): JobDefinition { return definition; } export interface WorkflowStep { name: string; run(input: I): O | Promise; } export function defineWorkflow(name: string, steps: Array>) { return { name, steps, async run(input: T): Promise { let value: unknown = input; for (const step of steps) value = await step.run(value); return value; }, }; } export function cronToInterval(cron: string): number { const aliases: Record = { "@hourly": 60 * 60 * 1000, "@daily": 24 * 60 * 60 * 1000, "@weekly": 7 * 24 * 60 * 60 * 1000, }; if (aliases[cron]) return aliases[cron]; const everyMinutes = /^\*\/(\d+)\s+\*\s+\*\s+\*\s+\*$/.exec(cron.trim()); if (everyMinutes) return Number(everyMinutes[1]) * 60 * 1000; throw new Error(`WRN-CRON-UNSUPPORTED: '${cron}'. Use @hourly, @daily, @weekly, or */N * * * *.`); } export function createQueue(options: QueueOptions = {}): Queue { const defaultMax = options.maxAttempts ?? 3; const backoffMs = options.backoffMs ?? 1000; const pollMs = options.pollMs ?? 250; const concurrency = options.concurrency ?? Number.POSITIVE_INFINITY; const capacity = options.capacity ?? Number.POSITIVE_INFINITY; if (!Number.isInteger(defaultMax) || defaultMax < 1) throw new RangeError("queue maxAttempts must be a positive integer"); if (!Number.isFinite(backoffMs) || backoffMs < 0) throw new RangeError("queue backoffMs must be a non-negative number"); if (!Number.isFinite(pollMs) || pollMs < 1) throw new RangeError("queue pollMs must be at least 1ms"); if (!( concurrency === Number.POSITIVE_INFINITY || (Number.isInteger(concurrency) && concurrency > 0) )) throw new RangeError("queue concurrency must be a positive integer"); if (!(capacity === Number.POSITIVE_INFINITY || (Number.isInteger(capacity) && capacity > 0))) throw new RangeError("queue capacity must be a positive integer"); const now = options.now ?? Date.now; const jobs: Job[] = []; const handlers = new Map(); const deadLetters = new Map(); const active = new Map }>(); let seq = 0; let timer: ReturnType | null = null; let draining = false; let accepting = true; async function runJob(job: Job): Promise { const handler = handlers.get(job.name); if (!handler) return; // no worker registered yet — leave it queued const idx = jobs.indexOf(job); if (idx >= 0) jobs.splice(idx, 1); // claim it job.attempts++; const controller = new AbortController(); const execution = (async () => { 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) { jobs.push({ ...job, attempts: 0, runAt: now() + job.repeat }); // recurring } } catch (error) { if (controller.signal.aborted) return; if (job.attempts < job.maxAttempts) { job.runAt = now() + backoffMs * Math.pow(2, job.attempts - 1); // exponential backoff jobs.push(job); } else { deadLetters.set(job.id, { ...job }); await options.onFailed?.(job, error); } } finally { active.delete(job.id); } })(); active.set(job.id, { controller, promise: execution }); await execution; } const drain: Queue["drain"] = async (at) => { if (draining) return 0; draining = true; try { const cutoff = at ?? now(); const due = jobs .filter((j) => j.runAt <= cutoff && handlers.has(j.name)) .sort((a, b) => b.priority - a.priority || a.runAt - b.runAt || a.createdAt - b.createdAt) .slice(0, concurrency); await Promise.all(due.map(runJob)); return due.length; } finally { draining = false; } }; return { async add(name, data, opts = {}) { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); if (!name.trim()) throw new TypeError("queue job name cannot be empty"); if ( opts.maxAttempts !== undefined && (!Number.isInteger(opts.maxAttempts) || opts.maxAttempts < 1) ) throw new RangeError("job maxAttempts must be a positive integer"); if (opts.delayMs !== undefined && (!Number.isFinite(opts.delayMs) || opts.delayMs < 0)) throw new RangeError("job delayMs must be a non-negative number"); if (opts.repeat !== undefined && (!Number.isFinite(opts.repeat) || opts.repeat <= 0)) throw new RangeError("job repeat must be a positive number"); if (opts.priority !== undefined && !Number.isFinite(opts.priority)) throw new RangeError("job priority must be a finite number"); if (opts.idempotencyKey) { const existing = jobs.find((job) => job.idempotencyKey === opts.idempotencyKey); if (existing) return existing as Job; } if (jobs.length + active.size >= capacity) throw new Error(`WRN-QUEUE-CAPACITY: queue capacity of ${capacity} reached`); const createdAt = now(); const job: Job = { id: `job_${++seq}`, name, data, attempts: 0, maxAttempts: opts.maxAttempts ?? defaultMax, runAt: createdAt + (opts.delayMs ?? 0), repeat: opts.repeat, priority: opts.priority ?? 0, idempotencyKey: opts.idempotencyKey, createdAt, }; jobs.push(job); return job as Job; }, process(name, handler) { handlers.set(name, handler as JobHandler); }, drain, start() { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); if (timer) return; timer = setInterval(() => void drain(), pollMs); }, stop() { if (timer) clearInterval(timer); timer = null; }, 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)); }, size: () => jobs.length, get: (id) => { const job = jobs.find((candidate) => candidate.id === id); return job ? { ...job } : undefined; }, list: (name) => jobs.filter((job) => !name || job.name === name).map((job) => ({ ...job })), cancel(id) { const index = jobs.findIndex((job) => job.id === id); if (index >= 0) { jobs.splice(index, 1); return true; } const running = active.get(id); if (!running) return false; running.controller.abort(); return true; }, failed: () => [...deadLetters.values()].map((job) => ({ ...job })), async retry(id) { if (!accepting) throw new Error("WRN-QUEUE-CLOSED: queue is shutting down"); const job = deadLetters.get(id); if (!job) return false; deadLetters.delete(id); jobs.push({ ...job, attempts: 0, runAt: now() }); return true; }, }; } export { memoryQueueStore, createDurableQueue } from "./durable.ts"; export type { QueueStore, DurableQueue, DurableQueueOptions } from "./durable.ts"; export { redisQueueStore, postgresQueueStore, POSTGRES_QUEUE_SCHEMA } from "./stores.ts"; export type { RedisQueueClient, SqlQueueClient } from "./stores.ts"; export { createQueueScheduler, addBatch, queueDashboardSnapshot, renderQueueDashboard, runQueueDaemon, } from "./scheduler.ts"; export type { ScheduledJob, QueueScheduler, QueueDashboardSnapshot } from "./scheduler.ts"; export { createWorkflowEngine, defineDurableWorkflow, memoryWorkflowStore } from "./workflow.ts"; export type { WorkflowDefinition, WorkflowEngine, WorkflowRunContext, WorkflowSnapshot, WorkflowStatus, WorkflowStore, } from "./workflow.ts"; export { subjectQueue } from "./subject.ts"; export type { SubjectJob, SubjectQueue } from "./subject.ts"; import { createExecutionContext, type ExecutionContext } from "@wrnexus/core";