/** Sequential job worker with retry and stale-job recovery. */ import { claimNextJob, updateJob } from "./db.js"; import { SkipJob } from "./review.js"; const STALE_MS = 2 * 60 * 60 * 1000; export class JobQueue { constructor({ db, engine, logger = console, pollMs = 3000, maxAttempts = 2 }) { this.db = db; this.engine = engine; this.logger = logger; this.pollMs = pollMs; this.maxAttempts = maxAttempts; this.running = false; this.busy = false; this.timer = null; this.currentJobId = null; } start() { if (this.running) return; this.running = true; this.recoverStale(); this.tick(); } stop() { this.running = false; if (this.timer) clearTimeout(this.timer); } recoverStale() { const cutoff = new Date(Date.now() - STALE_MS).toISOString().replace("T", " ").slice(0, 19); const stale = this.db.prepare( "SELECT id FROM jobs WHERE status = 'running' AND started_at IS NOT NULL AND started_at < ?", ).all(cutoff); for (const row of stale) { updateJob(this.db, row.id, { status: "queued", phase: null, error: "requeued after worker restart or timeout", }); this.logger.warn?.(`requeued stale job #${row.id}`); } } async tick() { if (!this.running) return; if (this.busy) { this.timer = setTimeout(() => this.tick(), this.pollMs); return; } const job = claimNextJob(this.db); if (!job) { this.timer = setTimeout(() => this.tick(), this.pollMs); return; } this.busy = true; this.currentJobId = job.id; try { await this.engine.execute(job); } catch (err) { if (err instanceof SkipJob) { await this.engine.failJob(job, err); } else if (job.attempts < this.maxAttempts) { this.logger.warn?.(`job #${job.id} failed (attempt ${job.attempts}): ${err.message}; retrying`); updateJob(this.db, job.id, { status: "queued", phase: null, error: `${err.name || "Error"}: ${err.message}`, }); } else { await this.engine.failJob(job, err); } } finally { this.busy = false; this.currentJobId = null; } this.timer = setTimeout(() => this.tick(), this.pollMs); } }