import { DatabaseSync } from 'node:sqlite'; import { mkdirSync } from 'node:fs'; import { join } from 'node:path'; import type { PrSnapshot, TimelineEvent, TimelineKind, TimelinePage } from '../shared/types.js'; import { rootDir, type MuteRule } from './config.js'; const SCHEMA = ` CREATE TABLE IF NOT EXISTS meta ( key TEXT PRIMARY KEY, value TEXT ); CREATE TABLE IF NOT EXISTS prs ( id TEXT PRIMARY KEY, repo TEXT NOT NULL, number INTEGER NOT NULL, state TEXT NOT NULL DEFAULT 'OPEN', signature TEXT NOT NULL, snapshot TEXT NOT NULL, detail_synced_at TEXT, updated_at TEXT NOT NULL, UNIQUE (repo, number) ); CREATE TABLE IF NOT EXISTS events ( id TEXT PRIMARY KEY, pr_id TEXT NOT NULL REFERENCES prs(id), kind TEXT NOT NULL, actor TEXT NOT NULL, created_at TEXT NOT NULL, updated_at TEXT, body_excerpt TEXT, url TEXT NOT NULL, review_state TEXT, thread_id TEXT, path TEXT, read INTEGER NOT NULL DEFAULT 0, read_at TEXT ); CREATE TABLE IF NOT EXISTS failure_reports ( pr_id TEXT NOT NULL, key TEXT NOT NULL, label TEXT NOT NULL, run_id TEXT NOT NULL, fetched_at TEXT NOT NULL, failed INTEGER NOT NULL, truncated INTEGER NOT NULL DEFAULT 0, failures TEXT NOT NULL, PRIMARY KEY (pr_id, key) ); CREATE INDEX IF NOT EXISTS idx_events_created ON events (created_at DESC); CREATE INDEX IF NOT EXISTS idx_events_pr ON events (pr_id); CREATE INDEX IF NOT EXISTS idx_events_unread ON events (read) WHERE read = 0; `; export interface FailureReportRow { key: string; label: string; fetched_at: string; failed: number; truncated: number; failures: string; } export interface EventInsert { id: string; prId: string; kind: TimelineKind; actor: string; createdAt: string; updatedAt: string | null; bodyExcerpt: string; url: string; reviewState: string | null; threadId: string | null; path: string | null; diffHunk: string | null; read: boolean; } export class Db { private db: DatabaseSync; /** Mute rules compiled to a SQL fragment; applied when *reading* the timeline and * unread counts, so config edits retroactively hide/unhide without data loss. */ private muteSql = ''; private muteParams: string[] = []; constructor(opts: { path?: string; muteRules?: MuteRule[] } = {}) { let path = opts.path; if (!path) { mkdirSync(join(rootDir, 'data'), { recursive: true }); path = join(rootDir, 'data', 'pr-monitor.sqlite'); } this.db = new DatabaseSync(path); this.db.exec('PRAGMA journal_mode = WAL'); this.db.exec('PRAGMA foreign_keys = ON'); this.db.exec(SCHEMA); const eventCols = this.db.prepare('PRAGMA table_info(events)').all() as { name: string }[]; if (!eventCols.some((c) => c.name === 'diff_hunk')) { this.db.exec('ALTER TABLE events ADD COLUMN diff_hunk TEXT'); } const clauses: string[] = []; for (const rule of opts.muteRules ?? []) { let clause = 'e.actor = ?'; this.muteParams.push(rule.actor); if (rule.kind) { clause += ' AND e.kind = ?'; this.muteParams.push(rule.kind); } if (rule.bodyContains) { clause += ' AND e.body_excerpt LIKE ?'; this.muteParams.push(`%${rule.bodyContains}%`); } clauses.push(`(${clause})`); } if (clauses.length > 0) this.muteSql = ` AND NOT (${clauses.join(' OR ')})`; } getMeta(key: string): string | null { const row = this.db.prepare('SELECT value FROM meta WHERE key = ?').get(key) as | { value: string | null } | undefined; return row?.value ?? null; } setMeta(key: string, value: string | null): void { this.db .prepare( 'INSERT INTO meta (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value = excluded.value', ) .run(key, value); } getPrSignatures(): Map { const rows = this.db .prepare("SELECT id, signature, detail_synced_at FROM prs WHERE state = 'OPEN'") .all() as { id: string; signature: string; detail_synced_at: string | null }[]; return new Map(rows.map((r) => [r.id, { signature: r.signature, detailSyncedAt: r.detail_synced_at }])); } upsertPr(snapshot: PrSnapshot, signature: string, detailSyncedAt: string | null): void { this.db .prepare( `INSERT INTO prs (id, repo, number, state, signature, snapshot, detail_synced_at, updated_at) VALUES (?, ?, ?, 'OPEN', ?, ?, ?, ?) ON CONFLICT(id) DO UPDATE SET state = 'OPEN', signature = excluded.signature, snapshot = excluded.snapshot, detail_synced_at = COALESCE(excluded.detail_synced_at, prs.detail_synced_at), updated_at = excluded.updated_at`, ) .run( snapshot.id, snapshot.repo, snapshot.number, signature, JSON.stringify(snapshot), detailSyncedAt, snapshot.updatedAt, ); } /** Update signature/snapshot from list data only (no detail fetch happened). */ markGone(ids: string[]): void { if (ids.length === 0) return; const stmt = this.db.prepare("UPDATE prs SET state = 'GONE' WHERE id = ?"); for (const id of ids) stmt.run(id); } getOpenSnapshots(): PrSnapshot[] { const rows = this.db .prepare("SELECT snapshot FROM prs WHERE state = 'OPEN' ORDER BY updated_at DESC") .all() as { snapshot: string }[]; return rows.map((r) => JSON.parse(r.snapshot) as PrSnapshot); } /** Insert or update an event. The conflict path never touches `read`. */ upsertEvent(e: EventInsert): void { this.db .prepare( `INSERT INTO events (id, pr_id, kind, actor, created_at, updated_at, body_excerpt, url, review_state, thread_id, path, diff_hunk, read, read_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, NULL) ON CONFLICT(id) DO UPDATE SET body_excerpt = excluded.body_excerpt, updated_at = excluded.updated_at, diff_hunk = excluded.diff_hunk`, ) .run( e.id, e.prId, e.kind, e.actor, e.createdAt, e.updatedAt, e.bodyExcerpt, e.url, e.reviewState, e.threadId, e.path, e.diffHunk, e.read ? 1 : 0, ); } getFailureReportRunId(prId: string, key: string): string | null { const row = this.db .prepare('SELECT run_id FROM failure_reports WHERE pr_id = ? AND key = ?') .get(prId, key) as { run_id: string } | undefined; return row?.run_id ?? null; } upsertFailureReport( prId: string, key: string, label: string, runId: string, report: { failed: number; truncated: boolean; failures: unknown[] }, ): void { this.db .prepare( `INSERT INTO failure_reports (pr_id, key, label, run_id, fetched_at, failed, truncated, failures) VALUES (?, ?, ?, ?, ?, ?, ?, ?) ON CONFLICT(pr_id, key) DO UPDATE SET label = excluded.label, run_id = excluded.run_id, fetched_at = excluded.fetched_at, failed = excluded.failed, truncated = excluded.truncated, failures = excluded.failures`, ) .run( prId, key, label, runId, new Date().toISOString(), report.failed, report.truncated ? 1 : 0, JSON.stringify(report.failures), ); } /** Remove reports whose rule no longer has failing checks. Returns true if any removed. */ pruneFailureReports(prId: string, activeKeys: string[]): boolean { const placeholders = activeKeys.map(() => '?').join(', '); const sql = activeKeys.length ? `DELETE FROM failure_reports WHERE pr_id = ? AND key NOT IN (${placeholders})` : 'DELETE FROM failure_reports WHERE pr_id = ?'; return Number(this.db.prepare(sql).run(prId, ...activeKeys).changes) > 0; } failureReportRows(prId: string): FailureReportRow[] { return this.db .prepare( 'SELECT key, label, fetched_at, failed, truncated, failures FROM failure_reports WHERE pr_id = ? ORDER BY label', ) .all(prId) as unknown as FailureReportRow[]; } failureSummaries(): Record { const rows = this.db .prepare( `SELECT f.pr_id, f.key, f.label, f.failed, f.truncated FROM failure_reports f JOIN prs p ON p.id = f.pr_id AND p.state = 'OPEN' ORDER BY f.label`, ) .all() as { pr_id: string; key: string; label: string; failed: number; truncated: number }[]; const out: Record = {}; for (const r of rows) { (out[r.pr_id] ??= []).push({ key: r.key, label: r.label, failed: r.failed, truncated: r.truncated === 1, }); } return out; } getSnapshot(id: string): PrSnapshot | null { const row = this.db.prepare('SELECT snapshot FROM prs WHERE id = ?').get(id) as | { snapshot: string } | undefined; return row ? (JSON.parse(row.snapshot) as PrSnapshot) : null; } hasEvent(id: string): boolean { return this.db.prepare('SELECT 1 FROM events WHERE id = ?').get(id) !== undefined; } timeline(opts: { limit: number; before?: string; unreadOnly?: boolean; prId?: string; }): TimelinePage { const clauses = ["p.state != ''"]; const params: (string | number)[] = []; if (opts.prId) { clauses.push('e.pr_id = ?'); params.push(opts.prId); } if (opts.before) { clauses.push('(e.created_at < ? OR (e.created_at = ? AND e.id < ?))'); const [createdAt, id] = splitCursor(opts.before); params.push(createdAt, createdAt, id); } if (opts.unreadOnly) clauses.push('e.read = 0'); const rows = this.db .prepare( `SELECT e.*, p.repo, p.number AS pr_number, p.snapshot, r.root_id, r.root_actor, r.root_body FROM events e JOIN prs p ON p.id = e.pr_id LEFT JOIN ( SELECT thread_id, id AS root_id, actor AS root_actor, body_excerpt AS root_body, MIN(created_at) FROM events WHERE thread_id IS NOT NULL GROUP BY thread_id ) r ON r.thread_id = e.thread_id WHERE ${clauses.join(' AND ')}${this.muteSql} ORDER BY e.created_at DESC, e.id DESC LIMIT ?`, ) .all(...params, ...this.muteParams, opts.limit + 1) as Record[]; const hasMore = rows.length > opts.limit; const page = rows.slice(0, opts.limit); const events: TimelineEvent[] = page.map((r) => ({ id: r.id as string, prId: r.pr_id as string, repo: r.repo as string, prNumber: r.pr_number as number, prTitle: (JSON.parse(r.snapshot as string) as { title: string }).title, kind: r.kind as TimelineEvent['kind'], actor: r.actor as string, createdAt: r.created_at as string, bodyExcerpt: (r.body_excerpt as string) ?? '', url: r.url as string, reviewState: r.review_state as string | null, path: r.path as string | null, read: (r.read as number) === 1, diffHunk: (r.diff_hunk as string | null) ?? null, inReplyTo: r.root_id && r.root_id !== r.id ? { actor: r.root_actor as string, bodyExcerpt: (r.root_body as string) ?? '' } : null, })); const last = events[events.length - 1]; return { events, nextCursor: hasMore && last ? joinCursor(last.createdAt, last.id) : null }; } markRead(id: string, read: boolean, at: string): void { this.db .prepare('UPDATE events SET read = ?, read_at = ? WHERE id = ?') .run(read ? 1 : 0, read ? at : null, id); } markAllRead(at: string, before?: string): number { const result = before ? this.db .prepare('UPDATE events SET read = 1, read_at = ? WHERE read = 0 AND created_at <= ?') .run(at, before) : this.db.prepare('UPDATE events SET read = 1, read_at = ? WHERE read = 0').run(at); return Number(result.changes); } markAllReadForPr(prId: string, at: string): number { const result = this.db .prepare('UPDATE events SET read = 1, read_at = ? WHERE read = 0 AND pr_id = ?') .run(at, prId); return Number(result.changes); } unreadCounts(): { byPr: Record; total: number } { const rows = this.db .prepare( `SELECT e.pr_id, COUNT(*) AS n FROM events e JOIN prs p ON p.id = e.pr_id AND p.state = 'OPEN' WHERE e.read = 0${this.muteSql} GROUP BY e.pr_id`, ) .all(...this.muteParams) as { pr_id: string; n: number }[]; const byPr: Record = {}; let total = 0; for (const r of rows) { byPr[r.pr_id] = r.n; total += r.n; } return { byPr, total }; } } export function joinCursor(createdAt: string, id: string): string { return `${createdAt}|${id}`; } function splitCursor(cursor: string): [string, string] { const idx = cursor.indexOf('|'); if (idx === -1) return [cursor, '']; return [cursor.slice(0, idx), cursor.slice(idx + 1)]; }