import { DatabaseSync } from "node:sqlite"; import { mkdirSync } from "node:fs"; import { dirname } from "node:path"; const SCHEMA = ` PRAGMA journal_mode = WAL; PRAGMA foreign_keys = ON; CREATE TABLE IF NOT EXISTS settings ( key TEXT PRIMARY KEY, value TEXT NOT NULL, updated_at TEXT NOT NULL DEFAULT (datetime('now')) ); CREATE TABLE IF NOT EXISTS repositories ( id INTEGER PRIMARY KEY AUTOINCREMENT, owner TEXT NOT NULL, name TEXT NOT NULL, repo_url TEXT, enabled INTEGER NOT NULL DEFAULT 1, -- The single branch every reviewed branch is merged into. managed_branch TEXT NOT NULL DEFAULT 'main', -- Comma-separated branches that are reviewed and may be merged up. check_branches TEXT NOT NULL DEFAULT 'main', review_scope TEXT NOT NULL DEFAULT 'both', gitea_token TEXT, llm_provider TEXT, llm_model TEXT, llm_base_url TEXT, llm_token TEXT, rule_path TEXT, background_template TEXT, excludes TEXT, publish_mode TEXT NOT NULL DEFAULT 'inline', create_issue INTEGER NOT NULL DEFAULT 1, issue_labels TEXT NOT NULL DEFAULT 'code-review', block_severity TEXT NOT NULL DEFAULT 'critical,high', block_categories TEXT NOT NULL DEFAULT '', fail_on_findings INTEGER NOT NULL DEFAULT 0, auto_merge INTEGER NOT NULL DEFAULT 0, auto_merge_mode TEXT NOT NULL DEFAULT 'when_checks_succeed', merge_method TEXT NOT NULL DEFAULT 'squash', delete_branch INTEGER NOT NULL DEFAULT 0, max_comments INTEGER NOT NULL DEFAULT 30, concurrency INTEGER NOT NULL DEFAULT 4, last_seen_at TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')), updated_at TEXT NOT NULL DEFAULT (datetime('now')), UNIQUE (owner, name) ); CREATE TABLE IF NOT EXISTS branch_state ( repo_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, ref_name TEXT NOT NULL, last_sha TEXT NOT NULL, reviewed_at TEXT NOT NULL DEFAULT (datetime('now')), PRIMARY KEY (repo_id, ref_name) ); CREATE TABLE IF NOT EXISTS jobs ( id INTEGER PRIMARY KEY AUTOINCREMENT, repo_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, trigger TEXT NOT NULL, ref_name TEXT NOT NULL, base_ref TEXT, from_sha TEXT, to_sha TEXT NOT NULL, pr_number INTEGER, status TEXT NOT NULL DEFAULT 'queued', phase TEXT, attempts INTEGER NOT NULL DEFAULT 0, findings INTEGER NOT NULL DEFAULT 0, blocking INTEGER NOT NULL DEFAULT 0, comment_count INTEGER NOT NULL DEFAULT 0, issue_number INTEGER, merged INTEGER NOT NULL DEFAULT 0, result_json TEXT, error TEXT, log TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')), started_at TEXT, finished_at TEXT ); CREATE TABLE IF NOT EXISTS review_records ( id INTEGER PRIMARY KEY AUTOINCREMENT, job_id INTEGER NOT NULL REFERENCES jobs(id) ON DELETE CASCADE, repo_id INTEGER NOT NULL REFERENCES repositories(id) ON DELETE CASCADE, ref_name TEXT NOT NULL, from_sha TEXT, to_sha TEXT NOT NULL, pr_number INTEGER, pr_url TEXT, trigger TEXT, source TEXT, requirement TEXT, params_json TEXT, findings_json TEXT, llm_model TEXT, llm_status TEXT, summary_status TEXT NOT NULL DEFAULT 'pending', summary TEXT, summary_error TEXT, summary_model TEXT, tokens_total INTEGER DEFAULT 0, elapsed TEXT, created_at TEXT NOT NULL DEFAULT (datetime('now')), summarized_at TEXT ); CREATE INDEX IF NOT EXISTS idx_review_records_repo ON review_records (repo_id, id DESC); CREATE INDEX IF NOT EXISTS idx_review_records_job ON review_records (job_id); CREATE INDEX IF NOT EXISTS idx_review_records_summary ON review_records (summary_status, id); CREATE INDEX IF NOT EXISTS idx_jobs_status ON jobs (status, id); CREATE INDEX IF NOT EXISTS idx_jobs_repo ON jobs (repo_id, id DESC); CREATE INDEX IF NOT EXISTS idx_jobs_sha ON jobs (repo_id, to_sha); CREATE TABLE IF NOT EXISTS webhook_deliveries ( delivery_id TEXT PRIMARY KEY, received_at TEXT NOT NULL DEFAULT (datetime('now')) ); `; function tableColumns(db, table) { try { return new Set(db.prepare(`PRAGMA table_info(${table})`).all().map((r) => r.name)); } catch { return new Set(); } } function addColumnIfMissing(db, table, column, definition) { if (tableColumns(db, table).has(column)) return false; db.exec(`ALTER TABLE ${table} ADD COLUMN ${column} ${definition}`); return true; } /** * Bring an existing database up to the current schema. * SQLite cannot drop or rename columns in place, so the pre-branch-model * `base_branch` / `branch_patterns` columns are left in place but ignored; * their values seed the new `managed_branch` / `check_branches` columns once. */ function migrate(db) { const addedManaged = addColumnIfMissing(db, "repositories", "managed_branch", "TEXT NOT NULL DEFAULT 'main'"); addColumnIfMissing(db, "repositories", "check_branches", "TEXT NOT NULL DEFAULT 'main'"); addColumnIfMissing(db, "repositories", "repo_url", "TEXT"); addColumnIfMissing(db, "jobs", "record_id", "INTEGER"); const cols = tableColumns(db, "repositories"); if (addedManaged && cols.has("base_branch")) { // Seed from the legacy single-branch configuration. db.exec("UPDATE repositories SET managed_branch = COALESCE(NULLIF(base_branch, ''), 'main')"); } if (cols.has("branch_patterns")) { // A legacy glob list becomes the explicit list of branches to check. const rows = db.prepare("SELECT id, branch_patterns, check_branches FROM repositories").all(); for (const row of rows) { const legacy = String(row.branch_patterns || "").trim(); const current = String(row.check_branches || "").trim(); if (!legacy || legacy === "*" || current !== "main") continue; const branches = legacy.split(",").map((x) => x.trim()).filter((x) => x && !x.includes("*") && !x.includes("?")); if (branches.length) { db.prepare("UPDATE repositories SET check_branches = ? WHERE id = ?").run(branches.join(","), row.id); } } } } export function openDatabase(path) { mkdirSync(dirname(path), { recursive: true }); const db = new DatabaseSync(path); db.exec(SCHEMA); migrate(db); return db; } export function getSetting(db, key, fallback = null) { const row = db.prepare("SELECT value FROM settings WHERE key = ?").get(key); return row ? row.value : fallback; } export function setSetting(db, key, value) { db.prepare( `INSERT INTO settings (key, value, updated_at) VALUES (?, ?, datetime('now')) ON CONFLICT (key) DO UPDATE SET value = excluded.value, updated_at = datetime('now')`, ).run(key, value == null ? "" : String(value)); } export function allSettings(db) { const out = {}; for (const row of db.prepare("SELECT key, value FROM settings").all()) { out[row.key] = row.value; } return out; } export function listRepositories(db) { return db.prepare("SELECT * FROM repositories ORDER BY owner, name").all(); } export function getRepository(db, id) { return db.prepare("SELECT * FROM repositories WHERE id = ?").get(id); } export function findRepository(db, owner, name) { return db.prepare("SELECT * FROM repositories WHERE owner = ? AND name = ?").get(owner, name); } const REPO_FIELDS = [ "owner", "name", "repo_url", "enabled", "managed_branch", "check_branches", "review_scope", "gitea_token", "llm_provider", "llm_model", "llm_base_url", "llm_token", "rule_path", "background_template", "excludes", "publish_mode", "create_issue", "issue_labels", "block_severity", "block_categories", "fail_on_findings", "auto_merge", "auto_merge_mode", "merge_method", "delete_branch", "max_comments", "concurrency", ]; export function upsertRepository(db, input) { const existing = input.id ? getRepository(db, input.id) : findRepository(db, input.owner, input.name); const values = {}; for (const field of REPO_FIELDS) { if (input[field] !== undefined) values[field] = input[field]; } if (existing) { const sets = Object.keys(values).map((k) => `${k} = ?`); if (sets.length > 0) { db.prepare( `UPDATE repositories SET ${sets.join(", ")}, updated_at = datetime('now') WHERE id = ?`, ).run(...Object.values(values), existing.id); } return getRepository(db, existing.id); } const cols = ["owner", "name", ...Object.keys(values)]; const params = [input.owner, input.name, ...Object.values(values)]; const placeholders = cols.map(() => "?").join(", "); const info = db.prepare( `INSERT INTO repositories (${cols.join(", ")}) VALUES (${placeholders})`, ).run(...params); return getRepository(db, Number(info.lastInsertRowid)); } export function deleteRepository(db, id) { db.prepare("DELETE FROM repositories WHERE id = ?").run(id); } export function enqueueJob(db, job) { const info = db.prepare( `INSERT INTO jobs (repo_id, trigger, ref_name, base_ref, from_sha, to_sha, pr_number, status) VALUES (?, ?, ?, ?, ?, ?, ?, 'queued')`, ).run( job.repoId, job.trigger, job.refName, job.baseRef ?? null, job.fromSha ?? null, job.toSha, job.prNumber ?? null, ); return Number(info.lastInsertRowid); } export function claimNextJob(db) { const row = db.prepare( "SELECT * FROM jobs WHERE status = 'queued' ORDER BY id LIMIT 1", ).get(); if (!row) return null; const info = db.prepare( `UPDATE jobs SET status = 'running', attempts = attempts + 1, started_at = datetime('now'), phase = 'starting' WHERE id = ? AND status = 'queued'`, ).run(row.id); if (info.changes === 0) return null; return db.prepare("SELECT * FROM jobs WHERE id = ?").get(row.id); } export function updateJob(db, id, patch) { const allowed = [ "status", "phase", "findings", "blocking", "comment_count", "issue_number", "merged", "result_json", "error", "log", "pr_number", ]; const keys = Object.keys(patch).filter((k) => allowed.includes(k)); if (keys.length === 0) return; const sets = keys.map((k) => `${k} = ?`); if (patch.status && ["succeeded", "failed", "skipped"].includes(patch.status)) { sets.push("finished_at = datetime('now')"); } db.prepare(`UPDATE jobs SET ${sets.join(", ")} WHERE id = ?`).run( ...keys.map((k) => patch[k]), id, ); } export function getJob(db, id) { return db.prepare("SELECT * FROM jobs WHERE id = ?").get(id); } export function listJobs(db, { repoId, limit = 50 } = {}) { if (repoId) { return db.prepare( "SELECT j.*, r.owner, r.name FROM jobs j JOIN repositories r ON r.id = j.repo_id " + "WHERE j.repo_id = ? ORDER BY j.id DESC LIMIT ?", ).all(repoId, limit); } return db.prepare( "SELECT j.*, r.owner, r.name FROM jobs j JOIN repositories r ON r.id = j.repo_id " + "ORDER BY j.id DESC LIMIT ?", ).all(limit); } export function findJobBySha(db, repoId, sha) { return db.prepare( "SELECT * FROM jobs WHERE repo_id = ? AND to_sha = ? AND status IN ('queued','running') ORDER BY id DESC LIMIT 1", ).get(repoId, sha); } export function getBranchState(db, repoId, refName) { return db.prepare( "SELECT * FROM branch_state WHERE repo_id = ? AND ref_name = ?", ).get(repoId, refName); } export function setBranchState(db, repoId, refName, sha) { db.prepare( `INSERT INTO branch_state (repo_id, ref_name, last_sha, reviewed_at) VALUES (?, ?, ?, datetime('now')) ON CONFLICT (repo_id, ref_name) DO UPDATE SET last_sha = excluded.last_sha, reviewed_at = datetime('now')`, ).run(repoId, refName, sha); } export function recordDelivery(db, deliveryId) { try { db.prepare("INSERT INTO webhook_deliveries (delivery_id) VALUES (?)").run(deliveryId); return true; } catch { return false; } } export function pruneDeliveries(db, keep = 2000) { db.prepare( `DELETE FROM webhook_deliveries WHERE delivery_id IN ( SELECT delivery_id FROM webhook_deliveries ORDER BY received_at DESC LIMIT -1 OFFSET ? )`, ).run(keep); } /* ------------------------------ review records ---------------------------- */ export function createReviewRecord(db, record) { const info = db.prepare( `INSERT INTO review_records (job_id, repo_id, ref_name, from_sha, to_sha, pr_number, pr_url, trigger, source, requirement, params_json, findings_json, llm_model, llm_status, tokens_total, elapsed, summary_status) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, 'pending')`, ).run( record.jobId, record.repoId, record.refName, record.fromSha ?? null, record.toSha, record.prNumber ?? null, record.prUrl ?? null, record.trigger ?? null, record.source ?? null, record.requirement ?? null, record.paramsJson ?? null, record.findingsJson ?? null, record.llmModel ?? null, record.llmStatus ?? null, record.tokensTotal ?? 0, record.elapsed ?? null, ); const id = Number(info.lastInsertRowid); db.prepare("UPDATE jobs SET record_id = ? WHERE id = ?").run(id, record.jobId); return id; } export function updateReviewRecord(db, id, patch) { const allowed = [ "pr_url", "requirement", "params_json", "findings_json", "llm_model", "llm_status", "summary_status", "summary", "summary_error", "summary_model", "tokens_total", "elapsed", ]; const keys = Object.keys(patch).filter((k) => allowed.includes(k)); if (keys.length === 0) return; const sets = keys.map((k) => `${k} = ?`); if (patch.summary_status && ["done", "failed", "skipped"].includes(patch.summary_status)) { sets.push("summarized_at = datetime('now')"); } db.prepare(`UPDATE review_records SET ${sets.join(", ")} WHERE id = ?`).run( ...keys.map((k) => patch[k]), id, ); } export function getReviewRecord(db, id) { return db.prepare("SELECT * FROM review_records WHERE id = ?").get(id); } export function getReviewRecordByJob(db, jobId) { return db.prepare("SELECT * FROM review_records WHERE job_id = ?").get(jobId); } export function listReviewRecords(db, { repoId, limit = 100 } = {}) { if (repoId) { return db.prepare( `SELECT rr.*, r.owner, r.name FROM review_records rr JOIN repositories r ON r.id = rr.repo_id WHERE rr.repo_id = ? ORDER BY rr.id DESC LIMIT ?`, ).all(repoId, limit); } return db.prepare( `SELECT rr.*, r.owner, r.name FROM review_records rr JOIN repositories r ON r.id = rr.repo_id ORDER BY rr.id DESC LIMIT ?`, ).all(limit); } export function claimPendingSummary(db) { const row = db.prepare( "SELECT * FROM review_records WHERE summary_status = 'pending' ORDER BY id LIMIT 1", ).get(); if (!row) return null; const info = db.prepare( "UPDATE review_records SET summary_status = 'running' WHERE id = ? AND summary_status = 'pending'", ).run(row.id); if (info.changes === 0) return null; return db.prepare("SELECT * FROM review_records WHERE id = ?").get(row.id); } export function requeueStaleSummaries(db) { const info = db.prepare( "UPDATE review_records SET summary_status = 'pending' WHERE summary_status = 'running'", ).run(); return info.changes; } export function jobStats(db) { const rows = db.prepare("SELECT status, COUNT(*) AS n FROM jobs GROUP BY status").all(); const out = { queued: 0, running: 0, succeeded: 0, failed: 0, skipped: 0 }; for (const r of rows) out[r.status] = r.n; return out; }