import { type Job, type JobProgress, type EnqueueParams, type StepContext, PermanentError } from './types'; import { readJob, writeJob, listAllJobs } from './storage'; import { getHandler } from './handler-registry'; import { sendMail } from 'emailer'; const activeLanes = new Map(); const PROGRESS_THROTTLE_MS = 1000; export async function enqueue(params: EnqueueParams): Promise { const handler = getHandler(params.type); if (!handler) throw new Error(`No handler registered for job type: ${params.type}`); const job: Job = { id: crypto.randomUUID(), lane: params.lane, type: params.type, userId: params.userId, status: 'queued', steps: handler.steps.map((s) => ({ name: s.name, status: 'pending' as const })), currentStep: 0, createdAt: Date.now(), meta: params.meta, notify: params.notify, }; await writeJob(job); console.log(`[queue] Enqueued job ${job.id} (${job.type}) in lane ${job.lane}`); kickLane(job.lane); return job; } export async function cancelJob(id: string): Promise { const job = await readJob(id); if (!job) return null; if (job.status === 'completed' || job.status === 'failed' || job.status === 'cancelled') return job; job.status = 'cancelled'; job.completedAt = Date.now(); for (const step of job.steps) { if (step.status === 'pending' || step.status === 'running') { step.status = 'failed'; step.error = 'Cancelled'; } } await writeJob(job); console.log(`[queue] Cancelled job ${job.id}`); return job; } export async function resumeInterruptedJobs() { const jobs = await listAllJobs(); const lanesToKick = new Set(); for (const job of jobs) { if (job.status === 'running') { job.status = 'queued'; job.startedAt = undefined; for (const step of job.steps) { if (step.status === 'running') { step.status = 'pending'; step.startedAt = undefined; } } await writeJob(job); console.log(`[queue] Reset interrupted job ${job.id} back to queued`); lanesToKick.add(job.lane); } else if (job.status === 'queued') { lanesToKick.add(job.lane); } } for (const lane of lanesToKick) { kickLane(lane); } } function kickLane(lane: string) { if (activeLanes.get(lane)) return; activeLanes.set(lane, true); processNextInLane(lane); } function scheduleRetry(lane: string, delayMs: number) { setTimeout(() => kickLane(lane), delayMs); } async function processNextInLane(lane: string) { try { const jobs = await listAllJobs(); const now = Date.now(); const next = jobs .filter((j) => j.lane === lane && j.status === 'queued' && (!j.retryAt || j.retryAt <= now)) .sort((a, b) => a.createdAt - b.createdAt)[0]; if (!next) { activeLanes.set(lane, false); return; } await runJob(next); } catch (err) { console.error(`[queue] Lane ${lane} processing error:`, err); } finally { const jobs = await listAllJobs(); const hasMore = jobs.some( (j) => j.lane === lane && j.status === 'queued' && (!j.retryAt || j.retryAt <= Date.now()), ); if (hasMore) { processNextInLane(lane); } else { activeLanes.set(lane, false); } } } async function runJob(job: Job) { const handler = getHandler(job.type); if (!handler) { job.status = 'failed'; job.error = `No handler for type: ${job.type}`; job.completedAt = Date.now(); await writeJob(job); return; } job.status = 'running'; job.startedAt = Date.now(); job.retryAt = undefined; await writeJob(job); const isRetry = (job.retries ?? 0) > 0; console.log( `[queue] ${isRetry ? 'Resuming' : 'Running'} job ${job.id} (${job.type})${isRetry ? ` retry ${job.retries}` : ''}`, ); const sharedMeta: Record = { ...(job.meta ?? {}) }; for (let i = 0; i < handler.steps.length; i++) { const fresh = await readJob(job.id); if (!fresh || fresh.status === 'cancelled') { console.log(`[queue] Job ${job.id} was cancelled, stopping`); return; } const handlerStep = handler.steps[i]!; const step = fresh.steps[i]!; // Skip already-completed steps on retry if (step.status === 'completed') continue; fresh.currentStep = i; step.status = 'running'; step.startedAt = Date.now(); await writeJob(fresh); console.log(`[queue] Job ${fresh.id} step ${i + 1}/${handler.steps.length}: "${handlerStep.name}"`); let lastProgressWrite = 0; let pendingProgress: JobProgress | null = null; const updateProgress = async (progress: JobProgress) => { step.progress = progress; const now = Date.now(); if (now - lastProgressWrite >= PROGRESS_THROTTLE_MS) { lastProgressWrite = now; pendingProgress = null; await writeJob(fresh); } else { pendingProgress = progress; } }; const ctx: StepContext = { job: fresh, step, updateProgress, meta: sharedMeta }; try { await handlerStep.run(ctx); if (pendingProgress) { step.progress = pendingProgress; } step.status = 'completed'; step.completedAt = Date.now(); await writeJob(fresh); } catch (err) { const errorMessage = err instanceof Error ? err.message : String(err); step.status = 'failed'; step.error = errorMessage; step.completedAt = Date.now(); // Check if handler supports retry (PermanentError skips retries) const isPermanent = err instanceof PermanentError; const retries = (fresh.retries ?? 0) + 1; if (!isPermanent && handler.retry && retries <= handler.retry.maxRetries) { // Schedule retry: reset failed step to pending, re-queue step.status = 'pending'; step.error = undefined; step.startedAt = undefined; step.completedAt = undefined; step.progress = undefined; fresh.status = 'queued'; fresh.error = undefined; fresh.completedAt = undefined; fresh.startedAt = undefined; fresh.retries = retries; fresh.retryAt = Date.now() + handler.retry.delayMs; await writeJob(fresh); console.log( `[queue] Job ${fresh.id} will retry (${retries}/${handler.retry.maxRetries}) in ${handler.retry.delayMs / 1000}s — failed at step "${step.name}": ${errorMessage}`, ); scheduleRetry(fresh.lane, handler.retry.delayMs); return; } fresh.status = 'failed'; fresh.error = `Step "${step.name}" failed: ${errorMessage}`; fresh.completedAt = Date.now(); fresh.meta = { ...fresh.meta, ...sharedMeta }; await writeJob(fresh); console.error(`[queue] Job ${fresh.id} failed at step "${step.name}":`, errorMessage); if (fresh.notify !== false) await notifyFailure(fresh); return; } } const final = await readJob(job.id); if (final && final.status === 'running') { final.status = 'completed'; final.completedAt = Date.now(); final.meta = { ...final.meta, ...sharedMeta }; await writeJob(final); console.log(`[queue] Job ${final.id} completed`); if (final.notify !== false) await notifyCompletion(final); } } async function notifyCompletion(job: Job) { try { await sendMail({ template: 'JobCompleted', subject: `Job completed: ${job.type}`, to: job.userId, data: { job }, }); } catch { // SMTP might not be configured — non-fatal } } async function notifyFailure(job: Job) { try { await sendMail({ template: 'JobFailed', subject: `Job failed: ${job.type}`, to: job.userId, data: { job }, }); } catch { // SMTP might not be configured — non-fatal } }