import { type Job, type JobProgress, type EnqueueParams, type StepContext, PermanentError } from '../queue/types'; import { readJob, writeJob, listAllJobs, ensureQueueDir } from '../queue/storage'; import { getHandler } from '../queue/handler-registry'; // Import handlers to register them import '../queue/handlers'; function formatDuration(ms: number): string { const s = Math.floor(ms / 1000); if (s < 60) return `${s}s`; const m = Math.floor(s / 60); const rem = s % 60; if (m < 60) return rem > 0 ? `${m}m${rem}s` : `${m}m`; const h = Math.floor(m / 60); const remM = m % 60; return remM > 0 ? `${h}h${remM}m` : `${h}h`; } const activeLanes = new Map(); const PROGRESS_THROTTLE_MS = 1000; export async function initQueue() { await ensureQueueDir(); await resumeInterruptedJobs(); console.log('[sidecar:queue] initialized'); } 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(`[sidecar: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(`[sidecar:queue] cancelled job ${job.id}`); return job; } 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(`[sidecar:queue] reset interrupted job ${job.id} back to queued`); lanesToKick.add(job.lane); } else if (job.status === 'queued') { // Clear retry delay on restart — no reason to wait after a sidecar restart if (job.retryAt) { job.retryAt = undefined; await writeJob(job); console.log(`[sidecar:queue] cleared retry delay for job ${job.id}`); } 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(`[sidecar: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; const startTime = Date.now(); console.log( `[sidecar: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(`[sidecar:queue] job ${job.id} was cancelled, stopping`); return; } const handlerStep = handler.steps[i]!; const step = fresh.steps[i]!; if (step.status === 'completed') continue; fresh.currentStep = i; step.status = 'running'; step.startedAt = Date.now(); await writeJob(fresh); 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); console.error(`[sidecar:queue] step "${step.name}" failed: ${errorMessage}`); step.status = 'failed'; step.error = errorMessage; step.completedAt = Date.now(); const isPermanent = err instanceof PermanentError; const retries = (fresh.retries ?? 0) + 1; if (!isPermanent && handler.retry && retries <= handler.retry.maxRetries) { 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( `[sidecar:queue] job ${fresh.id} will retry (${retries}/${handler.retry.maxRetries}) in ${handler.retry.delayMs / 1000}s`, ); 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(`[sidecar:queue] ✗ job ${fresh.id} failed at step "${step.name}" in ${formatDuration(Date.now() - startTime)}:`, errorMessage); 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(`[sidecar:queue] ✓ job ${final.id} completed in ${formatDuration(Date.now() - startTime)}`); await notifyCompletion(final); } } async function notifyCompletion(job: Job) { try { const { sendMail } = await import('emailer'); 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 { const { sendMail } = await import('emailer'); await sendMail({ template: 'JobFailed', subject: `Job failed: ${job.type}`, to: job.userId, data: { job }, }); } catch { // SMTP might not be configured — non-fatal } } export { readJob, listAllJobs };