仓库接入 - 新增「从 Gitea 导入」:用全局 Token 列出可见仓库,一键建配置 - 新增仓库只填 Git 地址,自动解析 owner/name,支持 https/ssh/scp 写法 - 新增「测试连接」按钮,保存前即可校验地址并拉取分支 分支模型 - 由单一 base_branch + glob 改为「一个管理分支 + 多个检查分支」 - 管理分支是唯一合并目标;检查分支全部纳入监控 - 支持从任意分支克隆创建管理分支 - 审查基准改为管理分支的 merge-base;仅目标为管理分支的 PR 才自动合并 - 旧库自动迁移:base_branch 播种 managed_branch,branch_patterns 展开为检查分支 审查历史 - 新增 review_records 表与「审查历史」页 - 记录触发来源、PR 链接、审查范围、阻断阈值、发布方式、自动合并设置、 排除路径、LLM 模型、token 消耗、耗时、需求背景与全部审查意见 - 详情弹窗一览,支持按仓库过滤 AI 摘要 - 新增独立摘要模块(app/lib/summary.js),与代码审查提示词分离 - 专用提示词输出固定四节、500 字内的中文记录:结论/范围/问题/要点 - 与代码审查共用全局 LLM 设置;摘要失败不影响审查与合并,可单条重跑 修复 - 摘要改用内置 fetch:运行镜像没有 curl,原先 spawn curl 必然 ENOENT - 去掉 blob:none 部分克隆并把凭据写入 .git/config: 惰性取 blob 不会带上 per-command extraHeader,私有库会报 could not read Username UI - 审查背景改为多行文本域(可滚动) - 分支改为可点选列表,管理分支高亮 - 仓库表格展示管理分支与检查分支
146 lines
5.1 KiB
JavaScript
146 lines
5.1 KiB
JavaScript
/** 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;
|
|
}
|
|
} |