| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905 |
- /**
- * Database Layer
- *
- * Handles SQLite database initialization and connection management.
- */
- import { SqliteDatabase, SqliteBackend, createDatabase } from './sqlite-adapter';
- import * as fs from 'fs';
- import * as path from 'path';
- import { SchemaVersion } from '../types';
- import { runMigrations, getCurrentVersion, CURRENT_SCHEMA_VERSION } from './migrations';
- import { getCodeGraphDir } from '../directory';
- export { SqliteDatabase, SqliteBackend } from './sqlite-adapter';
- /**
- * Apply connection-level PRAGMAs. Shared by `initialize` and `open` so the two
- * paths can't drift.
- *
- * `busy_timeout` is set FIRST, before any pragma that might touch the database
- * file (notably `journal_mode`). If another process holds a write lock at open
- * time, the later pragmas — and the connection's first query — then wait out
- * the lock instead of throwing "database is locked" immediately. See issue #238.
- *
- * The 5s window (was 120s) rides out a normal incremental sync; the old
- * 2-minute wait presented as a frozen, hung agent. With WAL, reads never block
- * on a writer, so this timeout only governs cross-process write contention
- * (e.g. the git-hook `codegraph sync` running while the MCP server writes).
- */
- function configureConnection(db: SqliteDatabase): void {
- db.pragma('busy_timeout = 5000'); // MUST be first — see above
- db.pragma('foreign_keys = ON');
- db.pragma('journal_mode = WAL'); // node:sqlite supports WAL on every platform
- db.pragma('synchronous = NORMAL'); // safe with WAL mode
- db.pragma('cache_size = -64000'); // 64 MB page cache
- db.pragma('temp_store = MEMORY'); // temp tables in memory
- db.pragma('mmap_size = 268435456'); // 256 MB memory-mapped I/O
- // Without a journal_size_limit the -wal file never shrinks below its
- // high-water mark while a connection lives: checkpoints fold frames back but
- // leave the file at full size, so one giant deferred-sync WAL stays giant
- // forever. With the limit set, any checkpoint that resets the WAL truncates
- // the file back down. Killed-process leftovers are handled separately by
- // healOversizedWal() at open. (#1431)
- db.pragma(`journal_size_limit = ${WAL_HEAL_THRESHOLD_BYTES}`);
- }
- /**
- * WAL size past which `healOversizedWal` (run at every `open`) checkpoints and
- * truncates the file, and to which `journal_size_limit` clips the WAL after any
- * resetting checkpoint. A SIGKILL'd process (the #850 liveness watchdog, OOM,
- * crash) can leave an arbitrarily large WAL behind — a whole deferred-sync
- * run's worth (#1248) — and before #1431 no later session ever shrank it: the
- * file just grew, killed session after killed session, until the disk filled
- * (25.6 GB observed). 64 MB is far above anything a healthy open ever sees
- * (a clean close deletes the WAL) yet small enough to cap the leak.
- * Override with `CODEGRAPH_WAL_HEAL_MB` (also feeds `journal_size_limit`).
- */
- export const WAL_HEAL_THRESHOLD_BYTES = resolveWalHealBytes(process.env.CODEGRAPH_WAL_HEAL_MB);
- /** Resolve the heal threshold from the env override (MB); invalid ⇒ 64 MB. */
- export function resolveWalHealBytes(envVal: string | undefined): number {
- if (envVal !== undefined && envVal !== '') {
- const n = Number(envVal);
- if (Number.isFinite(n) && n > 0) return Math.floor(n * 1024 * 1024);
- }
- return 64 * 1024 * 1024;
- }
- /**
- * Database connection wrapper with lifecycle management
- */
- export class DatabaseConnection {
- private db: SqliteDatabase;
- private dbPath: string;
- private backend: SqliteBackend;
- /**
- * `dev:ino` of the DB file at the moment we opened it (or null when the
- * platform/filesystem reports no usable inode). Lets us notice when the file
- * we hold open has been unlinked and REPLACED by a new file at the same path
- * — a git worktree removed and re-added, or `.codegraph/` deleted and
- * re-`init`ed under a long-lived server — at which point our fd reads a now
- * dead inode forever (#925). See `isReplacedOnDisk`.
- */
- private openedInode: string | null;
- /**
- * Whether FTS5 is available in this Node.js build. When false, search
- * falls back to LIKE + fuzzy matching (#1532).
- */
- readonly fts5Available: boolean;
- private constructor(db: SqliteDatabase, dbPath: string, backend: SqliteBackend, fts5Available: boolean) {
- this.db = db;
- this.dbPath = dbPath;
- this.backend = backend;
- this.fts5Available = fts5Available;
- this.openedInode = statInode(dbPath);
- }
- /**
- * Initialize a new database at the given path
- */
- static initialize(dbPath: string): DatabaseConnection {
- // Ensure parent directory exists
- const dir = path.dirname(dbPath);
- if (!fs.existsSync(dir)) {
- fs.mkdirSync(dir, { recursive: true });
- }
- // Create and configure database
- const { db, backend } = createDatabase(dbPath);
- configureConnection(db);
- // Run schema initialization, splitting FTS5 from the rest so
- // codegraph still works when Node.js was built without FTS5 (#1532).
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- const FTS5_MARKER = '-- Full-text search index on node names, docstrings, and signatures';
- const ftsIdx = schema.indexOf(FTS5_MARKER);
- let fts5Available = true;
- if (ftsIdx >= 0) {
- const preFts = schema.slice(0, ftsIdx);
- // FTS ends after the update trigger; required tables and indexes follow
- // it in schema.sql and must still be created when FTS5 is unavailable.
- const ftsSection = schema.slice(ftsIdx).match(
- /^[\s\S]*?CREATE TRIGGER IF NOT EXISTS nodes_au\b[\s\S]*?END;/
- )?.[0];
- if (!ftsSection) throw new Error('schema.sql: FTS5 update trigger not found');
- // Execute everything before FTS5 first
- db.exec(preFts);
- // Try FTS5; if it fails, skip it and continue with LIKE-only search
- try {
- db.exec(ftsSection);
- } catch (err: any) {
- fts5Available = false;
- const msg = err?.message ?? String(err);
- console.warn(
- `[codegraph] FTS5 not available in this Node.js build (${msg}). ` +
- `Search will fall back to LIKE + fuzzy matching. ` +
- `For full-text search, use a Node.js build with FTS5 enabled.`
- );
- }
- db.exec(schema.slice(ftsIdx + ftsSection.length));
- } else {
- db.exec(schema);
- }
- // Record current schema version so migrations aren't re-applied on open
- const currentVersion = getCurrentVersion(db);
- if (currentVersion < CURRENT_SCHEMA_VERSION) {
- db.prepare(
- 'INSERT OR IGNORE INTO schema_versions (version, applied_at, description) VALUES (?, ?, ?)'
- ).run(CURRENT_SCHEMA_VERSION, Date.now(), 'Initial schema includes all migrations');
- }
- return new DatabaseConnection(db, dbPath, backend, fts5Available);
- }
- /**
- * Open an existing database
- */
- static open(dbPath: string): DatabaseConnection {
- if (!fs.existsSync(dbPath)) {
- throw new Error(`Database not found: ${dbPath}`);
- }
- const { db, backend } = createDatabase(dbPath);
- configureConnection(db);
- // Detect FTS5 availability for search fallback (#1532)
- let fts5Available = true;
- try {
- db.exec("SELECT * FROM nodes_fts LIMIT 0");
- } catch {
- fts5Available = false;
- }
- // Check and run migrations if needed
- const conn = new DatabaseConnection(db, dbPath, backend, fts5Available);
- const currentVersion = getCurrentVersion(db);
- if (currentVersion < CURRENT_SCHEMA_VERSION) {
- runMigrations(db, currentVersion);
- }
- // Self-heal a bulk-load window that never closed (crash between
- // beginBulkNodeLoad and endBulkNodeLoad): the FTS triggers are missing and
- // nodes_fts is stale. Rebuild + recreate so search stays in sync.
- conn.healBulkNodeLoad();
- conn.healBulkSecondaryIndexes();
- // Self-heal a killed session's leftover oversized WAL (#1431) — one
- // statSync when healthy, off-thread checkpoint+truncate when not.
- void conn.healOversizedWal();
- return conn;
- }
- /**
- * FTS maintenance triggers dropped/recreated around a bulk load.
- * Names must match schema.sql.
- */
- private static readonly FTS_TRIGGER_NAMES = ['nodes_ai', 'nodes_ad', 'nodes_au'] as const;
- /**
- * Enter bulk-load mode: drop the per-row FTS sync triggers so mass node
- * inserts skip per-row tokenization. MUST be paired with endBulkNodeLoad()
- * (use try/finally); a crash inside the window is healed on the next open().
- * The window is DB-wide (triggers are schema objects), which is safe because
- * endBulkNodeLoad() rebuilds nodes_fts from the nodes table wholesale — any
- * row written by anyone during the window is captured by the rebuild.
- */
- beginBulkNodeLoad(): void {
- if (!this.fts5Available) return;
- for (const t of DatabaseConnection.FTS_TRIGGER_NAMES) {
- this.db.exec(`DROP TRIGGER IF EXISTS ${t}`);
- }
- }
- /**
- * Leave bulk-load mode: rebuild the whole FTS index from the nodes table in
- * one pass (far cheaper than per-row trigger firings), then recreate the
- * triggers by re-running schema.sql (idempotent — everything in it is
- * IF NOT EXISTS).
- */
- endBulkNodeLoad(): void {
- if (!this.fts5Available) return;
- this.db.exec(`INSERT INTO nodes_fts(nodes_fts) VALUES('rebuild')`);
- this.recreateFtsTriggers();
- }
- /**
- * NON-UNIQUE secondary indexes maintained per-row during the parse phase's
- * bulk inserts — the store-architecture arc's first lever (plan §4d: dubbo's
- * parse-loop wall is 94% store-writer busy, and the #1320 post-mortem showed
- * statement batching and sorted inserts are ~zero on this path because
- * B-TREE MAINTENANCE is the floor). A fresh init writes every row of
- * nodes/unresolved_refs/files exactly once and reads none of them until
- * resolution, so the parse window can drop all of these and rebuild each in
- * one table scan afterwards — the same measured trade as the resolution
- * phase's edge-index window (2.8s → 1.1s inserting, ~0.3s recreating).
- * Primary keys and UNIQUE constraints stay (upserts and OR-IGNORE dedup
- * conflict on them).
- */
- private static readonly BULK_PARSE_INDEX_NAMES = [
- 'idx_nodes_kind',
- 'idx_nodes_name',
- 'idx_nodes_qualified_name',
- 'idx_nodes_file_path',
- 'idx_nodes_language',
- 'idx_nodes_file_line',
- 'idx_nodes_lower_name',
- 'idx_unresolved_from_node',
- 'idx_unresolved_name',
- 'idx_unresolved_file_path',
- 'idx_unresolved_from_name',
- 'idx_unresolved_status',
- 'idx_unresolved_failed_tail',
- 'idx_files_language',
- 'idx_files_modified_at',
- ] as const;
- /**
- * Enter bulk-parse-load mode (FRESH-INIT ONLY — the caller gates on a fresh
- * DB, because an incremental index deletes per-file rows mid-phase and needs
- * the file_path indexes): drop every parse-lane secondary index, including
- * the four non-unique edge indexes (parse inserts contains-edges too; the
- * UNIQUE identity index stays for INSERT OR IGNORE dedup, and its `source`
- * prefix keeps source-keyed reads indexed, as in the edge window). MUST be
- * paired with endBulkParseLoad(); a crash inside the window is healed on the
- * next DatabaseConnection open (schema.sql re-applies CREATE INDEX IF NOT
- * EXISTS).
- */
- beginBulkParseLoad(): void {
- for (const idx of DatabaseConnection.BULK_PARSE_INDEX_NAMES) {
- this.db.exec(`DROP INDEX IF EXISTS ${idx}`);
- }
- this.beginBulkEdgeLoad();
- }
- /**
- * Leave bulk-parse-load mode: recreate everything the window dropped, one
- * table scan per index, with a yield between statements (same
- * liveness-watchdog rationale as endBulkEdgeLoad — at kernel scale each
- * build is a long synchronous scan). The edge indexes are rebuilt here too,
- * so paths that never enter the resolution phase's own bulk-edge window
- * (small runs) are left with a complete schema; the batched resolver's
- * beginBulkEdgeLoad simply re-drops them (DROP IF EXISTS — idempotent).
- */
- async endBulkParseLoad(): Promise<void> {
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- for (const idx of DatabaseConnection.BULK_PARSE_INDEX_NAMES) {
- const m = schema.match(new RegExp(`CREATE INDEX IF NOT EXISTS ${idx}\\b[^;]*;`));
- if (!m) throw new Error(`schema.sql: parse index ${idx} not found for bulk-load recreation`);
- this.db.exec(m[0]);
- await new Promise((resolve) => setImmediate(resolve));
- }
- await this.endBulkEdgeLoad();
- }
- /**
- * unresolved_refs secondary indexes NOT read by the batched resolution
- * loop. The loop pages pending refs by keyset (`status='pending' AND id>?`
- * — the status index + PK), deletes resolved rows by id, and parks failures
- * with a status UPDATE; every other ref index serves SYNC-time paths
- * (per-file re-index deletes, name-keyed retry, failed-tail heal). Each
- * per-batch DELETE maintains all of them — the biggest single main-thread
- * stage on the dubbo profile (deletes 1.2s of a 5.4s resolution phase) —
- * so the batched loop drops them and rebuilds at the end, where the table
- * holds only the surviving FAILED refs (resolved rows are gone), making
- * the recreate near-free.
- */
- private static readonly BULK_REF_INDEX_NAMES = [
- 'idx_unresolved_from_node',
- 'idx_unresolved_name',
- 'idx_unresolved_file_path',
- 'idx_unresolved_from_name',
- 'idx_unresolved_failed_tail',
- ] as const;
- /**
- * Enter bulk-ref mode for the batched resolution loop — see
- * BULK_REF_INDEX_NAMES. MUST be paired with endBulkRefLoad(); a crash
- * inside the window heals on the next open (schema.sql re-applies
- * CREATE INDEX IF NOT EXISTS).
- */
- beginBulkRefLoad(): void {
- for (const idx of DatabaseConnection.BULK_REF_INDEX_NAMES) {
- this.db.exec(`DROP INDEX IF EXISTS ${idx}`);
- }
- }
- /** Leave bulk-ref mode: recreate each index in one scan (yield between). */
- async endBulkRefLoad(): Promise<void> {
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- for (const idx of DatabaseConnection.BULK_REF_INDEX_NAMES) {
- const m = schema.match(new RegExp(`CREATE INDEX IF NOT EXISTS ${idx}\\b[^;]*;`));
- if (!m) throw new Error(`schema.sql: ref index ${idx} not found for bulk-load recreation`);
- this.db.exec(m[0]);
- await new Promise((resolve) => setImmediate(resolve));
- }
- }
- /**
- * Names of the NON-UNIQUE edge indexes dropped for a bulk edge load.
- * idx_edges_identity deliberately stays: INSERT OR IGNORE's dedup conflicts
- * on it (#1034), and its leftmost column is `source`, so the source-keyed
- * reads resolution makes mid-window (supertype walks over
- * `implements`/`extends`) keep an index via its prefix — verified with
- * EXPLAIN QUERY PLAN. Target-keyed and kind-keyed reads (traversal,
- * synthesis) happen only after endBulkEdgeLoad().
- */
- private static readonly BULK_EDGE_INDEX_NAMES = [
- 'idx_edges_kind',
- 'idx_edges_source_kind',
- 'idx_edges_target_kind',
- 'idx_edges_provenance',
- ] as const;
- /**
- * Enter bulk-edge-load mode: drop the non-unique edge indexes so the mass
- * INSERT OR IGNORE stream pays one B-tree (the identity index) instead of
- * five — measured 2.8s → 1.1s inserting a 224k-edge resolution set, with
- * recreation costing ~0.3s. MUST be paired with endBulkEdgeLoad(); a crash
- * inside the window is healed on the next DatabaseConnection open (schema.sql
- * re-applies CREATE INDEX IF NOT EXISTS).
- */
- beginBulkEdgeLoad(): void {
- for (const idx of DatabaseConnection.BULK_EDGE_INDEX_NAMES) {
- this.db.exec(`DROP INDEX IF EXISTS ${idx}`);
- }
- }
- /**
- * Leave bulk-edge-load mode: recreate the dropped indexes in one pass each
- * over the (now fully loaded) edges table — far cheaper than maintaining
- * them per-insert. DDL is extracted from schema.sql so it cannot drift.
- *
- * Async with a yield BETWEEN the four CREATE INDEX statements: each build is
- * a synchronous scan of the whole edges table (~20s apiece at Linux-kernel
- * scale, 79s total measured), and running them back-to-back is a single
- * event-loop stall longer than the #850 liveness watchdog's 60s window — a
- * daemon-triggered re-index would be SIGKILLed right after doing the work.
- * One yield per statement keeps every stall to a single index build, which
- * stays inside the window.
- */
- async endBulkEdgeLoad(): Promise<void> {
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- for (const idx of DatabaseConnection.BULK_EDGE_INDEX_NAMES) {
- const m = schema.match(new RegExp(`CREATE INDEX IF NOT EXISTS ${idx}\\b[^;]*;`));
- if (!m) throw new Error(`schema.sql: edge index ${idx} not found for bulk-load recreation`);
- this.db.exec(m[0]);
- await new Promise((resolve) => setImmediate(resolve));
- }
- }
- /** Recreate the FTS triggers + rebuild if a bulk-load window never closed. */
- private healBulkNodeLoad(): void {
- if (!this.fts5Available) return;
- const row = this.db
- .prepare(
- `SELECT count(*) AS c FROM sqlite_master WHERE type = 'trigger' AND name IN ('nodes_ai','nodes_ad','nodes_au')`
- )
- .get() as { c: number } | undefined;
- if ((row?.c ?? 0) >= DatabaseConnection.FTS_TRIGGER_NAMES.length) return;
- this.endBulkNodeLoad();
- }
- /** Recreate every secondary index a killed bulk parse/ref/edge window may leave dropped. */
- private healBulkSecondaryIndexes(): void {
- const names = [...new Set<string>([
- ...DatabaseConnection.BULK_PARSE_INDEX_NAMES,
- ...DatabaseConnection.BULK_REF_INDEX_NAMES,
- ...DatabaseConnection.BULK_EDGE_INDEX_NAMES,
- ])];
- const placeholders = names.map(() => '?').join(',');
- const row = this.db
- .prepare(`SELECT count(*) AS c FROM sqlite_master WHERE type = 'index' AND name IN (${placeholders})`)
- .get(...names) as { c: number } | undefined;
- if ((row?.c ?? 0) >= names.length) return;
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- for (const idx of names) {
- const m = schema.match(new RegExp(`CREATE INDEX IF NOT EXISTS ${idx}\\b[^;]*;`));
- if (!m) throw new Error(`schema.sql: index ${idx} not found for crash recovery`);
- this.db.exec(m[0]);
- }
- }
- /**
- * Recreate the FTS sync triggers from schema.sql — extracted from the file
- * rather than duplicated here so the DDL cannot drift from the schema.
- * (Re-execing the whole schema is not an option: it contains data INSERTs
- * that are not idempotent, e.g. schema_versions.)
- */
- private recreateFtsTriggers(): void {
- const schemaPath = path.join(__dirname, 'schema.sql');
- const schema = fs.readFileSync(schemaPath, 'utf-8');
- const triggerDdls = schema.match(
- /CREATE TRIGGER IF NOT EXISTS nodes_a[idu]\b[\s\S]*?END;/g
- );
- if (!triggerDdls || triggerDdls.length !== DatabaseConnection.FTS_TRIGGER_NAMES.length) {
- throw new Error(
- `schema.sql: expected ${DatabaseConnection.FTS_TRIGGER_NAMES.length} nodes FTS triggers, found ${triggerDdls?.length ?? 0}`
- );
- }
- for (const ddl of triggerDdls) {
- this.db.exec(ddl);
- }
- }
- /**
- * Get the underlying database instance
- */
- getDb(): SqliteDatabase {
- return this.db;
- }
- /**
- * Get the SQLite backend serving this connection. Per-instance so
- * MCP cross-project queries report the right backend even when
- * multiple project DBs are open in the same process.
- */
- getBackend(): SqliteBackend {
- return this.backend;
- }
- /**
- * Get database file path
- */
- getPath(): string {
- return this.dbPath;
- }
- /**
- * The journal mode actually in effect (e.g. 'wal', 'delete').
- *
- * SQLite silently keeps the prior mode if WAL can't be enabled — e.g. on
- * filesystems without shared-memory support (some network/virtualized mounts,
- * WSL2 /mnt). So the effective mode can differ
- * from what `configureConnection` requested. Surfaced in `codegraph status` so
- * a "database is locked" report is triageable: 'wal' ⇒ readers never block on a
- * writer; anything else ⇒ they can. See issue #238.
- */
- getJournalMode(): string {
- const raw = this.db.pragma('journal_mode');
- const row = Array.isArray(raw) ? raw[0] : raw;
- const mode = row && typeof row === 'object'
- ? (row as Record<string, unknown>).journal_mode
- : row;
- return String(mode ?? '').toLowerCase();
- }
- /**
- * Get current schema version
- */
- getSchemaVersion(): SchemaVersion | null {
- const row = this.db
- .prepare('SELECT version, applied_at, description FROM schema_versions ORDER BY version DESC LIMIT 1')
- .get() as { version: number; applied_at: number; description: string | null } | undefined;
- if (!row) return null;
- return {
- version: row.version,
- appliedAt: row.applied_at,
- description: row.description ?? undefined,
- };
- }
- /**
- * Execute a function within a transaction
- */
- transaction<T>(fn: () => T): T {
- return this.db.transaction(fn)();
- }
- /**
- * Get database file size in bytes
- */
- getSize(): number {
- const stats = fs.statSync(this.dbPath);
- return stats.size;
- }
- /**
- * Size of the `-wal` sidecar file in bytes. 0 when it doesn't exist (non-WAL
- * journal mode, in-memory DB, or no write since the last checkpoint+reset).
- */
- getWalSizeBytes(): number {
- if (!this.dbPath || this.dbPath === ':memory:') return 0;
- try {
- return fs.statSync(`${this.dbPath}-wal`).size;
- } catch {
- return 0;
- }
- }
- /** Size of the main DB file in bytes (0 for in-memory/unknown) — the WAL
- * valve scales its fold caps with it (resolveWalValveMb). */
- getDbFileSizeBytes(): number {
- if (!this.dbPath || this.dbPath === ':memory:') return 0;
- try {
- return fs.statSync(this.dbPath).size;
- } catch {
- return 0;
- }
- }
- /** Current `wal_autocheckpoint` interval in pages (0 = disabled). */
- getWalAutocheckpoint(): number {
- const v = this.db.pragma('wal_autocheckpoint', { simple: true });
- const n = Number(v);
- return Number.isFinite(n) ? n : 0;
- }
- /**
- * Set the connection's `wal_autocheckpoint` interval (pages; 0 disables).
- * Bulk indexing defers checkpoints entirely (#1231): the default 1000-page
- * auto-checkpoint re-writes hot B-tree/FTS pages into the main DB file over
- * and over — measured at ~95% of ALL disk I/O during a bulk index, and the
- * difference between 45s and 19+ minutes on HDD-class storage. During
- * deferral a {@link WalCheckpointValve} bounds WAL growth off-thread.
- */
- setWalAutocheckpoint(pages: number): void {
- this.db.pragma(`wal_autocheckpoint = ${Math.max(0, Math.floor(pages))}`);
- }
- /**
- * `PRAGMA wal_checkpoint(PASSIVE)` on a worker thread with its own
- * connection. PASSIVE never blocks the writer, and running it off-thread
- * means the main thread — and the #850 watchdog heartbeat — keep turning
- * even when the backfill is minutes of I/O on slow storage (a synchronous
- * checkpoint that exceeds the watchdog's 60s window gets a healthy index
- * SIGKILLed — observed in the #1231 repro).
- *
- * Returns SQLite's checkpoint result row — `log === checkpointed` with
- * `busy === 0` means the ENTIRE WAL was backfilled, so the writer's next
- * commit restarts the WAL from the top and the file stops growing. The
- * WAL valve needs that signal because a WAL file's SIZE never shrinks:
- * after the first wrap, raw file size says nothing about the un-backfilled
- * backlog. Best-effort: returns null on any failure (including worker
- * threads being unavailable — a potentially minutes-long checkpoint must
- * never run inline on the main thread).
- */
- async checkpointWalPassive(): Promise<{ busy: number; log: number; checkpointed: number } | null> {
- return this.checkpointWal('PASSIVE');
- }
- /**
- * `PRAGMA wal_checkpoint(TRUNCATE)` — same off-thread pattern as PASSIVE,
- * but on success the WAL FILE is chopped to zero. A completed passive
- * backfill bounds the un-checkpointed backlog, yet the FILE only stops
- * growing when a commit finds ZERO readers holding WAL marks — rare while
- * pool workers cycle, so at kernel scale a fully-backfilled WAL still
- * accreted the phase's whole write volume on disk (§7a.1: 22GB). The valve
- * calls this exactly at a parked barrier (writer parked, pool drained,
- * backfill complete) where the no-reader condition is guaranteed rather
- * than lucky. The worker sets a short busy_timeout so a racing reader
- * degrades this to a no-op (busy=1) instead of a stall.
- */
- async checkpointWalTruncate(): Promise<{ busy: number; log: number; checkpointed: number } | null> {
- return this.checkpointWal('TRUNCATE');
- }
- /**
- * Shrink a leftover oversized WAL (#1431). A SIGKILL'd session — the #850
- * liveness watchdog, OOM, a crash — leaves its WAL on disk, the next session
- * appends to the same file, and (pre-#1431) nothing ever truncated it:
- * PASSIVE checkpoints fold frames but keep the file at its high-water mark,
- * and the one shrinking path (a clean last-connection close) is exactly what
- * the killed world never takes. Unbounded growth until the disk fills.
- *
- * Called fire-and-forget from every `open()`: cost is one statSync when the
- * WAL is small (the overwhelmingly common case). Past the threshold it runs
- * the off-thread PASSIVE fold then TRUNCATE — both on worker connections
- * with a busy_timeout, so a racing writer degrades this to a no-op that the
- * next open retries rather than a stall.
- */
- async healOversizedWal(): Promise<{ healed: boolean; beforeBytes: number; afterBytes: number }> {
- const beforeBytes = this.getWalSizeBytes();
- if (beforeBytes <= WAL_HEAL_THRESHOLD_BYTES) {
- return { healed: false, beforeBytes, afterBytes: beforeBytes };
- }
- // Single-flight: open() fires this fire-and-forget and callers may also
- // invoke it explicitly. Two concurrent passes DEFEAT each other — each
- // checkpoint worker sees the other as a busy reader and no-ops — so share
- // one in-flight pass instead of racing.
- this.walHeal ??= this.runWalHeal(beforeBytes).finally(() => { this.walHeal = null; });
- return this.walHeal;
- }
- private walHeal: Promise<{ healed: boolean; beforeBytes: number; afterBytes: number }> | null = null;
- private async runWalHeal(beforeBytes: number): Promise<{ healed: boolean; beforeBytes: number; afterBytes: number }> {
- // A racing reader/writer (another session healing the same file, a query
- // pool warming up) degrades a checkpoint pass to a busy no-op — retry a
- // few times before leaving the rest to the next open.
- for (let attempt = 0; attempt < 3; attempt++) {
- if (attempt > 0) await new Promise((r) => setTimeout(r, 300));
- await this.checkpointWalPassive();
- await this.checkpointWalTruncate();
- if (this.getWalSizeBytes() <= WAL_HEAL_THRESHOLD_BYTES) break;
- }
- const afterBytes = this.getWalSizeBytes();
- if (process.env.CODEGRAPH_WAL_VALVE_DEBUG) {
- console.error(`[wal-heal] oversized WAL at open: ${Math.round(beforeBytes / (1024 * 1024))}MB -> ${Math.round(afterBytes / (1024 * 1024))}MB`);
- }
- return { healed: afterBytes < beforeBytes, beforeBytes, afterBytes };
- }
- private async checkpointWal(mode: 'PASSIVE' | 'TRUNCATE'): Promise<{ busy: number; log: number; checkpointed: number } | null> {
- if (!this.dbPath || this.dbPath === ':memory:') {
- try {
- const row = this.db.prepare(`PRAGMA wal_checkpoint(${mode})`).get() as Record<string, number> | undefined;
- return row ? { busy: Number(row.busy), log: Number(row.log), checkpointed: Number(row.checkpointed) } : null;
- } catch {
- return null;
- }
- }
- try {
- const { Worker } = await import('node:worker_threads');
- const workerSource = `
- const { workerData, parentPort } = require('node:worker_threads');
- let row = null;
- let err = null;
- try {
- const { DatabaseSync } = require('node:sqlite');
- const db = new DatabaseSync(workerData.dbPath);
- const mode = workerData.mode === 'TRUNCATE' ? 'TRUNCATE' : 'PASSIVE';
- try {
- if (mode === 'TRUNCATE') db.exec('PRAGMA busy_timeout = 2000');
- row = db.prepare('PRAGMA wal_checkpoint(' + mode + ')').get();
- } catch (e) { err = String(e && e.message || e); }
- try { db.close(); } catch {}
- } catch (e) { err = err || String(e && e.message || e); }
- parentPort.postMessage({ row, err });
- `;
- return await new Promise((resolve) => {
- let settled = false;
- const finish = (row?: Record<string, number> | null): void => {
- if (settled) return;
- settled = true;
- resolve(row ? { busy: Number(row.busy), log: Number(row.log), checkpointed: Number(row.checkpointed) } : null);
- };
- try {
- const worker = new Worker(workerSource, { eval: true, workerData: { dbPath: this.dbPath, mode } });
- worker.once('message', (m: { row?: Record<string, number> | null; err?: string | null }) => {
- if (m?.err && process.env.CODEGRAPH_WAL_VALVE_DEBUG) {
- console.error(`[wal-valve] checkpoint worker (${mode}): ${m.err}`);
- }
- void worker.terminate();
- finish(m?.row ?? null);
- });
- worker.once('error', () => { void worker.terminate(); finish(null); });
- worker.once('exit', () => finish(null));
- } catch {
- finish(null);
- }
- });
- } catch {
- return null;
- }
- }
- /**
- * Optimize database (vacuum and analyze)
- */
- optimize(): void {
- this.db.exec('VACUUM');
- this.db.exec('ANALYZE');
- }
- /**
- * Lightweight maintenance to run after bulk writes (indexAll, sync).
- * Two operations:
- *
- * - `PRAGMA optimize` — incremental ANALYZE; SQLite only re-analyzes
- * tables whose row counts changed materially since the last
- * ANALYZE. Without it, the query planner has no statistics on the
- * freshly-bulk-loaded tables and can pick suboptimal indexes.
- *
- * - `PRAGMA wal_checkpoint(PASSIVE)` — fold pending WAL pages back
- * into the main database file so the WAL file doesn't grow
- * unboundedly between automatic checkpoints (auto-fires at 1000
- * pages by default; large indexAll runs blow past that).
- *
- * Runs on a WORKER THREAD with its own connection: on a multi-GB index
- * these pragmas are minutes of synchronous IO (a 95k-file kernel index
- * left a 593MB WAL whose checkpoint alone blew the #850 watchdog's 60s
- * window and got a COMPLETED index SIGKILLed at the finish line). WAL
- * checkpointing from a second connection is standard SQLite; `PRAGMA
- * optimize` persists its statistics in sqlite_stat tables, so the main
- * connection benefits the same. The main thread just awaits a message,
- * so the event loop — and the watchdog heartbeat — keep turning.
- *
- * Everything is silently swallowed on failure — best-effort
- * optimization, never load-bearing for correctness. If worker threads
- * are unavailable, falls back to a bounded in-line `PRAGMA optimize`
- * and SKIPS the checkpoint (the final close() checkpoints after the
- * CLI has already disarmed its watchdog).
- */
- async runMaintenance(): Promise<void> {
- // In-memory / test databases: nothing worth a worker round-trip.
- if (!this.dbPath || this.dbPath === ':memory:') {
- try { this.db.exec('PRAGMA optimize'); } catch { /* ignore */ }
- try { this.db.exec('PRAGMA wal_checkpoint(PASSIVE)'); } catch { /* ignore */ }
- return;
- }
- await this.runPragmasOffThread(
- ['PRAGMA analysis_limit=1000', 'PRAGMA optimize', 'PRAGMA wal_checkpoint(PASSIVE)'],
- // Worker threads unavailable — bounded in-line fallback, no checkpoint.
- ['PRAGMA analysis_limit=1000', 'PRAGMA optimize']
- );
- }
- /**
- * Run pragmas on a worker thread against its own connection to this DB
- * (shared machinery for {@link runMaintenance} and
- * {@link checkpointWalPassive}). Each pragma is individually best-effort;
- * the whole call is best-effort. `inlineFallback` (if any) runs on THIS
- * connection only when worker threads are unavailable — keep it to pragmas
- * that are safe to run synchronously on the main thread.
- */
- private async runPragmasOffThread(pragmas: string[], inlineFallback: string[] = []): Promise<void> {
- try {
- const { Worker } = await import('node:worker_threads');
- const workerSource = `
- const { workerData, parentPort } = require('node:worker_threads');
- try {
- const { DatabaseSync } = require('node:sqlite');
- const db = new DatabaseSync(workerData.dbPath);
- for (const p of workerData.pragmas) { try { db.exec(p); } catch {} }
- try { db.close(); } catch {}
- } catch {}
- parentPort.postMessage('done');
- `;
- await new Promise<void>((resolve) => {
- let settled = false;
- const finish = (): void => {
- if (!settled) { settled = true; resolve(); }
- };
- try {
- const worker = new Worker(workerSource, { eval: true, workerData: { dbPath: this.dbPath, pragmas } });
- worker.once('message', () => { void worker.terminate(); finish(); });
- worker.once('error', () => { void worker.terminate(); finish(); });
- worker.once('exit', finish);
- } catch {
- finish();
- }
- });
- } catch {
- for (const p of inlineFallback) {
- try { this.db.exec(p); } catch { /* ignore */ }
- }
- }
- }
- /**
- * Close the database connection
- */
- close(): void {
- this.db.close();
- }
- /**
- * Check if the database connection is open
- */
- isOpen(): boolean {
- return this.db.open;
- }
- /**
- * True when the DB file at our path has been REPLACED on disk since we opened
- * it — a different inode now lives at the same path, so the fd we still hold
- * points at a now-unlinked inode that can never receive new writes (#925).
- * The trigger is removing and recreating `.codegraph/` at the same path under
- * a long-lived process (`git worktree remove` + re-add, or `rm -rf
- * .codegraph` + `codegraph init`). Returns false when the inode is unchanged,
- * when the file is momentarily absent (mid-recreate — nothing to reopen onto
- * yet), or when the platform doesn't report a usable inode (Windows can't
- * unlink an open file and its st_ino is unreliable, so this never fires there).
- */
- isReplacedOnDisk(): boolean {
- if (this.openedInode === null) return false;
- const current = statInode(this.dbPath);
- return current !== null && current !== this.openedInode;
- }
- }
- /**
- * `dev:ino` for a path, or null if it can't be stat'd or the platform doesn't
- * report a usable inode. Windows st_ino is unreliable across handle reopens, so
- * we deliberately return null there — the deleted-but-open-inode hazard this
- * guards (#925) is a POSIX file-semantics issue that doesn't arise on Windows
- * (an open file can't be unlinked).
- */
- function statInode(p: string): string | null {
- if (process.platform === 'win32') return null;
- try {
- const s = fs.statSync(p);
- return `${s.dev}:${s.ino}`;
- } catch {
- return null;
- }
- }
- /**
- * Default database filename
- */
- export const DATABASE_FILENAME = 'codegraph.db';
- /**
- * SQLite's sidecar files in WAL mode — the write-ahead log and its shared-memory
- * index. They sit beside the main DB file and are removed alongside it when the
- * database is discarded (see `removeDatabaseFiles`).
- */
- const WAL_SIDECAR_SUFFIXES = ['-wal', '-shm'] as const;
- /**
- * Get the default database path for a project
- */
- export function getDatabasePath(projectRoot: string): string {
- return path.join(getCodeGraphDir(projectRoot), DATABASE_FILENAME);
- }
- /**
- * Delete a database file and its WAL sidecars (`-wal`/`-shm`).
- *
- * This is how a FULL re-index discards an existing database — rather than
- * opening the old graph and DELETE-ing every row. On a large or pre-fix
- * poisoned index (e.g. an old graph that scanned an ignored gitlink corpus into
- * ~1.6M nodes with a multi-GB WAL, #1065) the per-row `nodes_fts` delete-trigger
- * churn blocks the main thread long enough to trip the #850 liveness watchdog
- * before indexing even starts, so the rebuild could never recover the bad state
- * (#1067). Unlinking is O(1) regardless of DB size and also reclaims the disk
- * the bloated WAL would otherwise keep.
- *
- * POSIX removes the directory entry even while another process (a daemon/MCP
- * server) still holds the file open; that holder heals via `reopenIfReplaced`
- * (#925). On Windows a live holder can make the unlink fail with EBUSY/EPERM —
- * that is thrown for the caller to surface ("stop the other process and retry").
- * The `-wal`/`-shm` sidecars are best-effort: SQLite recreates them on the next
- * open, so a leftover sidecar is harmless.
- */
- export function removeDatabaseFiles(dbPath: string): void {
- // The main DB file first — its removal is the operation that must succeed (or
- // report why it couldn't). force:true treats an already-missing file as done.
- fs.rmSync(dbPath, { force: true });
- for (const suffix of WAL_SIDECAR_SUFFIXES) {
- try {
- fs.rmSync(dbPath + suffix, { force: true });
- } catch {
- // A sidecar still held/locked is harmless — SQLite rebuilds it on open.
- }
- }
- }
|