import { DatabaseSync } from 'node:sqlite'; import { mkdirSync } from 'node:fs'; import { dirname } from 'node:path'; const DDL = [ `CREATE TABLE IF NOT EXISTS reviews ( id INTEGER PRIMARY KEY AUTOINCREMENT, repo TEXT NOT NULL, sha TEXT NOT NULL, pr INTEGER, event TEXT NOT NULL, pusher TEXT, authors TEXT, trailers TEXT, axes TEXT, overall INTEGER, grade TEXT, verdict TEXT, findings TEXT, status_state TEXT, error TEXT, diff_bytes INTEGER, truncated INTEGER DEFAULT 0, model TEXT, duration_ms INTEGER, ts TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ','now')) )`, `CREATE INDEX IF NOT EXISTS idx_reviews_repo_ts ON reviews(repo, ts DESC)`, `CREATE INDEX IF NOT EXISTS idx_reviews_sha ON reviews(sha)`, `CREATE TABLE IF NOT EXISTS queue_jobs ( id INTEGER PRIMARY KEY AUTOINCREMENT, repo TEXT NOT NULL, sha TEXT NOT NULL, job TEXT NOT NULL, state TEXT NOT NULL DEFAULT 'queued', outcome TEXT, created_ts TEXT NOT NULL DEFAULT (strftime('%Y-%m-%dT%H:%M:%fZ','now')), started_ts TEXT, finished_ts TEXT )`, `CREATE INDEX IF NOT EXISTS idx_queue_jobs_state ON queue_jobs(state)` ]; export function openDb(dbPath) { mkdirSync(dirname(dbPath), { recursive: true }); const db = new DatabaseSync(dbPath); for (const stmt of DDL) db.prepare(stmt).run(); ensureColumn(db, 'queue_jobs', 'started_ts', 'TEXT'); return db; } function ensureColumn(db, table, column, type) { const columns = db.prepare(`PRAGMA table_info(${table})`).all(); if (!columns.some((c) => c.name === column)) { db.prepare(`ALTER TABLE ${table} ADD COLUMN ${column} ${type}`).run(); } } export function insertReview(db, r) { const stmt = db.prepare(` INSERT INTO reviews (repo, sha, pr, event, pusher, authors, trailers, axes, overall, grade, verdict, findings, status_state, error, diff_bytes, truncated, model, duration_ms) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) `); const res = stmt.run( r.repo, r.sha, r.pr ?? null, r.event, r.pusher ?? null, JSON.stringify(r.authors ?? []), JSON.stringify(r.trailers ?? []), r.axes ? JSON.stringify(r.axes) : null, r.overall ?? null, r.grade ?? null, r.verdict ?? null, r.findings ? JSON.stringify(r.findings) : null, r.statusState ?? null, r.error ?? null, r.diffBytes ?? null, r.truncated ? 1 : 0, r.model ?? null, r.durationMs ?? null ); return res.lastInsertRowid; } // --- persistent work queue (crash-safe: rows written before the pending // status posts; rows that never reach 'done' are re-enqueued on startup) --- export function enqueueJob(db, job) { const res = db.prepare( `INSERT INTO queue_jobs (repo, sha, job) VALUES (?, ?, ?)` ).run(`${job.owner}/${job.repo}`, String(job.sha), JSON.stringify(job)); return res.lastInsertRowid; } export function startJob(db, id) { db.prepare( `UPDATE queue_jobs SET state = 'running', started_ts = strftime('%Y-%m-%dT%H:%M:%fZ','now') WHERE id = ?` ).run(id); } export function finishJob(db, id, outcome = null) { db.prepare( `UPDATE queue_jobs SET state = 'done', outcome = ?, finished_ts = strftime('%Y-%m-%dT%H:%M:%fZ','now') WHERE id = ?` ).run(outcome, id); } // Startup reconciliation: rows still 'queued' belonged to a previous process // that died mid-review (their posted status is stuck at 'pending'). export function recoverPendingJobs(db) { const rows = db.prepare( `SELECT id, job FROM queue_jobs WHERE state != 'done' ORDER BY id ASC` ).all(); const out = []; for (const row of rows) { try { out.push({ id: row.id, job: JSON.parse(row.job) }); } catch { finishJob(db, row.id, 'unparseable-job'); } } return out; } export function listReviews(db, repo, limit = 50) { return db.prepare( 'SELECT * FROM reviews WHERE repo = ? ORDER BY id DESC LIMIT ?' ).all(repo, limit); } export function allReviews(db, limit = 500) { return db.prepare('SELECT * FROM reviews ORDER BY id DESC LIMIT ?').all(limit); }