/** Sequential job worker with retry and stale-job recovery. */ import { claimNextJob, claimPendingSummary, requeueStaleSummaries, updateJob, updateReviewRecord, } from "./db.js"; import { SkipJob } from "./review.js"; import { summariseReview } from "./summary.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(); const requeued = requeueStaleSummaries(this.db); if (requeued) this.logger.warn?.(`requeued ${requeued} interrupted summariser task(s)`); 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) { // Nothing to review: use the idle slot for pending summaries. await this.runSummaryOnce(); 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; } // Summarising runs on the same worker so it never competes with a review // for the LLM endpoint. await this.runSummaryOnce(); this.timer = setTimeout(() => this.tick(), this.pollMs); } /** * Summarise one finished review, if any is pending. * Failures are recorded on the record and never block the queue. */ async runSummaryOnce() { const record = claimPendingSummary(this.db); if (!record) return false; try { const repo = this.db.prepare("SELECT * FROM repositories WHERE id = ?").get(record.repo_id); if (!repo) { updateReviewRecord(this.db, record.id, { summary_status: "skipped", summary_error: "repository no longer exists", }); return false; } let findings = []; try { findings = JSON.parse(record.findings_json || "[]"); } catch { findings = []; } const job = this.db.prepare("SELECT result_json FROM jobs WHERE id = ?").get(record.job_id); let summary = null; try { summary = JSON.parse(job?.result_json || "{}").summary ?? null; } catch { summary = null; } const llm = { url: repo.llm_base_url || this.engine.config.llmUrl, token: repo.llm_token || this.engine.config.llmToken, model: repo.llm_model || this.engine.config.llmModel, protocol: repo.llm_provider || this.engine.config.llmProtocol, authHeader: this.engine.config.llmAuthHeader, extraHeaders: this.engine.config.llmExtraHeaders, }; this.logger.info?.(`summarising review record #${record.id} (${repo.owner}/${repo.name} ${record.ref_name})`); const res = await summariseReview({ llm, repo, record, findings, summary, requirement: record.requirement, }); if (res.ok) { updateReviewRecord(this.db, record.id, { summary_status: "done", summary: res.text, summary_model: res.model, summary_error: res.complete ? null : "摘要缺少标准小节,已按原文保存", }); this.logger.info?.(`review record #${record.id} summarised (${res.text.length} chars)`); } else { updateReviewRecord(this.db, record.id, { summary_status: "failed", summary_error: String(res.error).slice(0, 500), }); this.logger.warn?.(`review record #${record.id} summary failed: ${res.error}`); } } catch (err) { updateReviewRecord(this.db, record.id, { summary_status: "failed", summary_error: String(err.message || err).slice(0, 500), }); this.logger.warn?.(`review record #${record.id} summary error: ${err.message}`); } return true; } }