diff --git a/gitnexus/src/core/embeddings/embedding-restore-spill.ts b/gitnexus/src/core/embeddings/embedding-restore-spill.ts new file mode 100644 index 000000000..10ca51b45 --- /dev/null +++ b/gitnexus/src/core/embeddings/embedding-restore-spill.ts @@ -0,0 +1,483 @@ +/** + * Disk-backed restore cache for CodeEmbedding rows (#3306). + * + * `loadCachedEmbeddings` used to `getAll()` the table and `map(Number)` every + * vector into a JS `number[]`. On a large already-indexed repo that single + * structure OOMs the V8 heap during "Caching embeddings..." even when the + * incremental diff is a handful of nodes. + * + * This module keeps metadata in RAM and writes vectors to a temp Float32 + * spill. Restore materializes only the rows that Phase 3.5 will re-insert, + * in the existing 200-row batches. + */ +import { closeSync, openSync, readSync, unlinkSync, writeSync } from 'node:fs'; +import { AsyncLocalStorage } from 'node:async_hooks'; +import { randomBytes } from 'node:crypto'; +import os from 'node:os'; +import path from 'node:path'; +import type { CachedEmbedding } from './types.js'; + +/** In-RAM vector copies stay below this row count; larger tables use the spill. */ +export const DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT = 2048; + +const SPILL_MAGIC = 'GNXE'; +const SPILL_VERSION = 1; +const SPILL_HEADER_BYTES = 12; + +export interface CachedEmbeddingMeta { + nodeId: string; + chunkIndex: number; + startLine: number; + endLine: number; + contentHash?: string; + /** Row order in the spill file (and in `embeddings` when in-memory). */ + vectorIndex: number; +} + +export interface EmbeddingVectorSpill { + path: string; + dims: number; + rowCount: number; +} + +export interface CachedEmbeddingsSnapshot { + embeddingNodeIds: Set; + /** Populated only when the table is at or under the in-memory row limit. */ + embeddings: CachedEmbedding[]; + rows: CachedEmbeddingMeta[]; + spill?: EmbeddingVectorSpill; +} + +export interface LoadCachedEmbeddingsOptions { + /** + * Keep full `number[]` vectors in RAM at or below this many rows. + * `0` always spills. Default {@link DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT} + * or `GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT`. + */ + inMemoryRowLimit?: number; + /** Directory for the spill file (default `os.tmpdir()`). */ + spillDir?: string; +} + +export interface CachedEmbeddingsBuilder { + embeddingNodeIds: Set; + rows: CachedEmbeddingMeta[]; + /** Float32 vectors kept in RAM until the in-memory row limit is exceeded. */ + inMemory: Float32Array[] | null; + inMemoryRowLimit: number; + writer: EmbeddingSpillWriter; +} + +export function emptyCachedEmbeddingsSnapshot(): CachedEmbeddingsSnapshot { + return { embeddingNodeIds: new Set(), embeddings: [], rows: [] }; +} + +export function resolveEmbeddingCacheInMemoryRowLimit(override?: number): number { + if (override !== undefined) { + if (!Number.isFinite(override) || override < 0) { + return DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT; + } + return Math.floor(override); + } + const raw = process.env.GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT; + if (raw === undefined || raw === '') return DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT; + const parsed = parseInt(raw, 10); + return Number.isFinite(parsed) && parsed >= 0 + ? parsed + : DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT; +} + +export function normalizeCachedEmbeddings(raw: { + embeddingNodeIds?: Set; + embeddings?: CachedEmbedding[]; + rows?: CachedEmbeddingMeta[]; + spill?: EmbeddingVectorSpill; +}): CachedEmbeddingsSnapshot { + const embeddings = raw.embeddings ?? []; + const embeddingNodeIds = raw.embeddingNodeIds ?? new Set(embeddings.map((row) => row.nodeId)); + const rows = + raw.rows ?? + embeddings.map((row, vectorIndex) => ({ + nodeId: row.nodeId, + chunkIndex: row.chunkIndex, + startLine: row.startLine, + endLine: row.endLine, + contentHash: row.contentHash, + vectorIndex, + })); + return { embeddingNodeIds, embeddings, rows, spill: raw.spill }; +} + +export function cacheRowCount(snapshot: CachedEmbeddingsSnapshot): number { + return snapshot.rows.length > 0 ? snapshot.rows.length : snapshot.embeddings.length; +} + +export function snapshotEmbeddingDims(snapshot: CachedEmbeddingsSnapshot): number | undefined { + if (snapshot.spill && snapshot.spill.dims > 0) return snapshot.spill.dims; + const dims = snapshot.embeddings[0]?.embedding.length; + return dims && dims > 0 ? dims : undefined; +} + +export function coerceEmbeddingToFloat32(embedding: unknown): Float32Array | null { + if (embedding == null) return null; + if (embedding instanceof Float32Array) { + return embedding.length > 0 ? embedding : null; + } + if (ArrayBuffer.isView(embedding) && !(embedding instanceof DataView)) { + const view = embedding as Exclude & { length: number }; + if (view.length === 0) return null; + return Float32Array.from({ length: view.length }, (_, i) => Number(view[i])); + } + if ( + typeof embedding === 'object' && + typeof (embedding as Iterable)[Symbol.iterator] === 'function' + ) { + const arr = Array.isArray(embedding) + ? (embedding as unknown[]) + : Array.from(embedding as Iterable); + if (arr.length === 0) return null; + return Float32Array.from(arr, (value) => Number(value)); + } + return null; +} + +export function float32ToNumberArray(vec: Float32Array): number[] { + const out = new Array(vec.length); + for (let i = 0; i < vec.length; i++) out[i] = vec[i]!; + return out; +} + +/** + * `fs.writeSync` can return a short byte count. Loop until the whole buffer + * lands, matching `sync-csv-writer.ts`, so a partial write never advances + * `rowCount` on a truncated vector. + */ +function writeAllSync(fd: number, data: Uint8Array): void { + let offset = 0; + while (offset < data.length) { + const n = writeSync(fd, data, offset, data.length - offset); + if (n <= 0) { + throw new Error(`embedding spill short write: wrote ${n} of ${data.length - offset} bytes`); + } + offset += n; + } +} + +function unlinkBestEffort(filePath: string): void { + try { + unlinkSync(filePath); + } catch { + /* ENOENT or already removed */ + } +} + +const liveSpillPaths = new Set(); +const spillScope = new AsyncLocalStorage>(); +let spillExitHookInstalled = false; + +function trackLiveSpillPath(filePath: string): void { + liveSpillPaths.add(filePath); + spillScope.getStore()?.add(filePath); + if (!spillExitHookInstalled) { + spillExitHookInstalled = true; + process.on('exit', () => { + for (const spillPath of liveSpillPaths) { + unlinkBestEffort(spillPath); + } + }); + } +} + +function untrackLiveSpillPath(filePath: string): void { + liveSpillPaths.delete(filePath); +} + +/** Best-effort unlink of every tracked spill. Safe to call more than once. */ +export function discardLiveEmbeddingSpills(): void { + for (const spillPath of [...liveSpillPaths]) { + unlinkBestEffort(spillPath); + liveSpillPaths.delete(spillPath); + } +} + +/** Run `fn` so later {@link discardScopedEmbeddingSpills} only unlinks this run. */ +export function withEmbeddingSpillScope(fn: () => T): T { + return spillScope.run(new Set(), fn); +} + +/** Unlink spills created inside the current {@link withEmbeddingSpillScope}. */ +export function discardScopedEmbeddingSpills(): void { + const owned = spillScope.getStore(); + if (!owned) return; + for (const spillPath of [...owned]) { + unlinkBestEffort(spillPath); + liveSpillPaths.delete(spillPath); + owned.delete(spillPath); + } +} + +export class EmbeddingSpillWriter { + readonly path: string; + dims = 0; + rowCount = 0; + private fd: number | null = null; + private closed = false; + + constructor(dir: string) { + this.path = path.join( + dir, + `gitnexus-embed-restore-${process.pid}-${randomBytes(8).toString('hex')}.bin`, + ); + } + + append(vec: Float32Array): void { + if (this.closed) { + throw new Error('embedding spill writer already closed'); + } + if (this.fd === null) { + this.dims = vec.length; + this.fd = openSync(this.path, 'wx', 0o600); + trackLiveSpillPath(this.path); + const header = Buffer.alloc(SPILL_HEADER_BYTES); + header.write(SPILL_MAGIC, 0, 4, 'ascii'); + header.writeUInt8(SPILL_VERSION, 4); + header.writeUInt32LE(this.dims, 5); + writeAllSync(this.fd, header); + } else if (vec.length !== this.dims) { + throw new Error( + `embedding dim mismatch while spilling: got ${vec.length}, expected ${this.dims}`, + ); + } + writeAllSync(this.fd, Buffer.from(vec.buffer, vec.byteOffset, vec.byteLength)); + this.rowCount++; + } + + finish(): EmbeddingVectorSpill | undefined { + if (this.fd !== null) { + closeSync(this.fd); + this.fd = null; + } + this.closed = true; + if (this.rowCount === 0) { + this.unlinkQuiet(); + return undefined; + } + return { path: this.path, dims: this.dims, rowCount: this.rowCount }; + } + + abort(): void { + const opened = this.fd !== null; + if (this.fd !== null) { + try { + closeSync(this.fd); + } catch { + /* already closed */ + } + this.fd = null; + } + this.closed = true; + if (opened || this.rowCount > 0) { + this.unlinkQuiet(); + } + } + + private unlinkQuiet(): void { + unlinkBestEffort(this.path); + untrackLiveSpillPath(this.path); + } +} + +/** Validates the spill header once and reads vectors without reopening the file. */ +export class EmbeddingSpillReader { + private fd: number | null = null; + private readonly bytesPerVec: number; + readonly dims: number; + readonly rowCount: number; + + constructor(spill: EmbeddingVectorSpill) { + this.rowCount = spill.rowCount; + const fd = openSync(spill.path, 'r'); + try { + const header = Buffer.alloc(SPILL_HEADER_BYTES); + const headerRead = readSync(fd, header, 0, SPILL_HEADER_BYTES, 0); + if (headerRead !== SPILL_HEADER_BYTES || header.toString('ascii', 0, 4) !== SPILL_MAGIC) { + throw new Error(`invalid embedding spill header: ${spill.path}`); + } + if (header.readUInt8(4) !== SPILL_VERSION) { + throw new Error(`unsupported embedding spill version in ${spill.path}`); + } + const dims = header.readUInt32LE(5); + if (dims !== spill.dims) { + throw new Error(`embedding spill dim mismatch: file ${dims}, expected ${spill.dims}`); + } + this.dims = dims; + this.bytesPerVec = dims * 4; + this.fd = fd; + } catch (err) { + closeSync(fd); + throw err; + } + } + + read(indices: readonly number[]): Float32Array[] { + if (this.fd === null) { + throw new Error('embedding spill reader already closed'); + } + const out: Float32Array[] = []; + for (const index of indices) { + if (!Number.isInteger(index) || index < 0 || index >= this.rowCount) { + throw new Error(`embedding spill index out of range: ${index}`); + } + const offset = SPILL_HEADER_BYTES + index * this.bytesPerVec; + const copy = new Float32Array(this.dims); + const bytes = new Uint8Array(copy.buffer, copy.byteOffset, this.bytesPerVec); + const n = readSync(this.fd, bytes, 0, this.bytesPerVec, offset); + if (n !== this.bytesPerVec) { + throw new Error(`short embedding spill read at index ${index}`); + } + out.push(copy); + } + return out; + } + + close(): void { + if (this.fd === null) return; + closeSync(this.fd); + this.fd = null; + } +} + +export function readSpillVectors( + spill: EmbeddingVectorSpill, + indices: readonly number[], +): Float32Array[] { + const reader = new EmbeddingSpillReader(spill); + try { + return reader.read(indices); + } finally { + reader.close(); + } +} + +export function disposeEmbeddingSpill(spill?: EmbeddingVectorSpill): void { + if (!spill?.path) return; + unlinkBestEffort(spill.path); + untrackLiveSpillPath(spill.path); +} + +export function createCachedEmbeddingsBuilder( + options?: LoadCachedEmbeddingsOptions, +): CachedEmbeddingsBuilder { + const inMemoryRowLimit = resolveEmbeddingCacheInMemoryRowLimit(options?.inMemoryRowLimit); + return { + embeddingNodeIds: new Set(), + rows: [], + inMemory: inMemoryRowLimit <= 0 ? null : [], + inMemoryRowLimit, + writer: new EmbeddingSpillWriter(options?.spillDir ?? os.tmpdir()), + }; +} + +export function ingestCachedEmbeddingRow( + builder: CachedEmbeddingsBuilder, + row: Record | unknown[], + hasContentHash: boolean, +): void { + const rec = row as Record & unknown[]; + const nodeId = String(rec.nodeId ?? rec[0] ?? ''); + if (!nodeId) return; + const embedding = rec.embedding ?? rec[4]; + const f32 = coerceEmbeddingToFloat32(embedding); + if (!f32) return; + + builder.embeddingNodeIds.add(nodeId); + const meta: CachedEmbeddingMeta = { + nodeId, + chunkIndex: Number(rec.chunkIndex ?? rec[1] ?? 0), + startLine: Number(rec.startLine ?? rec[2] ?? 0), + endLine: Number(rec.endLine ?? rec[3] ?? 0), + contentHash: hasContentHash + ? ((rec.contentHash ?? rec[5] ?? undefined) as string | undefined) + : undefined, + vectorIndex: builder.rows.length, + }; + builder.rows.push(meta); + + if (builder.inMemory && builder.rows.length <= builder.inMemoryRowLimit) { + builder.inMemory.push(f32); + return; + } + + if (builder.inMemory) { + for (const prior of builder.inMemory) { + builder.writer.append(prior); + } + builder.inMemory = null; + } + builder.writer.append(f32); +} + +export function finalizeCachedEmbeddingsSnapshot( + builder: CachedEmbeddingsBuilder, +): CachedEmbeddingsSnapshot { + const inMemory = builder.inMemory; + if (inMemory) { + builder.writer.abort(); + return { + embeddingNodeIds: builder.embeddingNodeIds, + embeddings: builder.rows.map((meta, i) => ({ + nodeId: meta.nodeId, + chunkIndex: meta.chunkIndex, + startLine: meta.startLine, + endLine: meta.endLine, + contentHash: meta.contentHash, + embedding: float32ToNumberArray(inMemory[i]!), + })), + rows: builder.rows, + }; + } + return { + embeddingNodeIds: builder.embeddingNodeIds, + embeddings: [], + rows: builder.rows, + spill: builder.writer.finish(), + }; +} + +export function abortCachedEmbeddingsBuilder(builder: CachedEmbeddingsBuilder): void { + builder.writer.abort(); +} + +export function materializeCachedEmbeddings( + snapshot: CachedEmbeddingsSnapshot, + metas: readonly CachedEmbeddingMeta[], + spillReader?: EmbeddingSpillReader, +): CachedEmbedding[] { + if (metas.length === 0) return []; + if (snapshot.spill && snapshot.embeddings.length === 0) { + const indices = metas.map((meta) => meta.vectorIndex); + const vectors = spillReader + ? spillReader.read(indices) + : readSpillVectors(snapshot.spill, indices); + return metas.map((meta, i) => ({ + nodeId: meta.nodeId, + chunkIndex: meta.chunkIndex, + startLine: meta.startLine, + endLine: meta.endLine, + contentHash: meta.contentHash, + embedding: float32ToNumberArray(vectors[i]!), + })); + } + if (snapshot.embeddings.length === 0) return []; + const byKey = new Map( + snapshot.embeddings.map((row) => [`${row.nodeId}:${row.chunkIndex}`, row] as const), + ); + return metas.map((meta) => { + const hit = + byKey.get(`${meta.nodeId}:${meta.chunkIndex}`) ?? snapshot.embeddings[meta.vectorIndex]; + if (!hit) { + throw new Error(`missing cached embedding ${meta.nodeId}:${meta.chunkIndex}`); + } + return hit; + }); +} diff --git a/gitnexus/src/core/lbug/lbug-adapter.ts b/gitnexus/src/core/lbug/lbug-adapter.ts index 621d26da6..5f5a076d4 100644 --- a/gitnexus/src/core/lbug/lbug-adapter.ts +++ b/gitnexus/src/core/lbug/lbug-adapter.ts @@ -34,7 +34,16 @@ import type { GraphEmitManifest } from './graph-emit-sink.js'; import type { PdgEmitManifest } from './pdg-emit-sink.js'; import { PDG_EDGE_TYPES } from './pdg-emit-sink.js'; import { getNodeLabel as deriveNodeLabel, type WriteStreamFactory } from './rel-pair-routing.js'; -import { EMBEDDABLE_LABELS, type CachedEmbedding } from '../embeddings/types.js'; +import { EMBEDDABLE_LABELS } from '../embeddings/types.js'; +import { + abortCachedEmbeddingsBuilder, + createCachedEmbeddingsBuilder, + emptyCachedEmbeddingsSnapshot, + finalizeCachedEmbeddingsSnapshot, + ingestCachedEmbeddingRow, + type CachedEmbeddingsSnapshot, + type LoadCachedEmbeddingsOptions, +} from '../embeddings/embedding-restore-spill.js'; import { extensionManager, getFtsCapability, @@ -2069,28 +2078,31 @@ export const getLbugStats = async (): Promise<{ /** * Load cached embeddings from LadybugDB before a rebuild. - * Returns all embedding vectors so they can be re-inserted after the graph is reloaded, - * avoiding expensive re-embedding of unchanged nodes. + * + * Streams `CodeEmbedding` rows with `hasNext`/`getNext` under `withConnLock` + * (#2264, #3306). Vectors are spilled to a temp Float32 file once the table + * exceeds the in-memory row limit so incremental analyze cannot OOM the V8 + * heap by materializing every `number[]` up front. Small tables still return + * in-RAM `embeddings` for existing callers/tests. * * Detects old schema (no chunkIndex column) and returns empty cache to trigger rebuild. */ -export const loadCachedEmbeddings = async (): Promise<{ - embeddingNodeIds: Set; - embeddings: CachedEmbedding[]; -}> => { +export const loadCachedEmbeddings = async ( + options?: LoadCachedEmbeddingsOptions, +): Promise => { const c = conn; if (!c) { - return { embeddingNodeIds: new Set(), embeddings: [] }; + return emptyCachedEmbeddingsSnapshot(); } // The whole read runs inside the connection lock (#2264 review P2). It's safe // today only by call-ordering (loadCachedEmbeddings runs before the WAL driver // starts), but the lock makes it robust to future reordering — a concurrent // CHECKPOINT on the singleton connection is the documented corruption trigger. - // Leaf read: no nested withConnLock-wrapped helpers inside. + // Leaf read: no nested withConnLock-wrapped helpers inside. Do NOT call + // `streamQuery` here — that path is unlocked and would race a CHECKPOINT. return withConnLock(async () => { - const embeddingNodeIds = new Set(); - const embeddings: CachedEmbedding[] = []; + const builder = createCachedEmbeddingsBuilder(options); try { // Schema migration detection: query with new columns to verify schema version. // Old schema only had (nodeId, embedding); new schema adds (id, chunkIndex, startLine, endLine, contentHash). @@ -2104,51 +2116,46 @@ export const loadCachedEmbeddings = async (): Promise<{ ); await readQueryRows(check); } catch { - return { embeddingNodeIds: new Set(), embeddings: [] }; + abortCachedEmbeddingsBuilder(builder); + return emptyCachedEmbeddingsSnapshot(); } - // Try to read contentHash alongside chunk columns - let rows: any; + let queryResult: lbug.QueryResult | lbug.QueryResult[] | undefined; let hasContentHash = true; try { - rows = await c.query( - `MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash`, - ); - } catch (err: any) { - // Fallback for legacy DBs without contentHash column - const msg = err?.message ?? ''; - if (isMissingColumnOrTableError(msg)) { - hasContentHash = false; - rows = await c.query( - `MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding`, + try { + queryResult = await c.query( + `MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash`, ); - } else { - throw err; + } catch (err: any) { + // Fallback for legacy DBs without contentHash column + const msg = err?.message ?? ''; + if (isMissingColumnOrTableError(msg)) { + hasContentHash = false; + queryResult = await c.query( + `MATCH (e:${EMBEDDING_TABLE_NAME}) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding`, + ); + } else { + throw err; + } } - } - for (const row of await readQueryRows(rows)) { - const nodeId = String(row.nodeId ?? row[0] ?? ''); - if (!nodeId) continue; - embeddingNodeIds.add(nodeId); - const embedding = row.embedding ?? row[4]; - if (embedding) { - embeddings.push({ - nodeId, - chunkIndex: Number(row.chunkIndex ?? row[1] ?? 0), - startLine: Number(row.startLine ?? row[2] ?? 0), - endLine: Number(row.endLine ?? row[3] ?? 0), - embedding: Array.isArray(embedding) - ? embedding.map(Number) - : Array.from(embedding as any).map(Number), - contentHash: hasContentHash ? (row.contentHash ?? row[5] ?? undefined) : undefined, - }); + const results = Array.isArray(queryResult) ? queryResult : [queryResult]; + const result = results[0]; + while (await result.hasNext()) { + const row = await result.getNext(); + ingestCachedEmbeddingRow(builder, row, hasContentHash); } + return finalizeCachedEmbeddingsSnapshot(builder); + } catch (err) { + abortCachedEmbeddingsBuilder(builder); + throw err; + } finally { + if (queryResult) await closeQueryResults(queryResult); } - } catch { - /* embedding table may not exist */ + } catch (err) { + abortCachedEmbeddingsBuilder(builder); + throw err; } - - return { embeddingNodeIds, embeddings }; }); }; diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index 10188637e..8ea0b0f66 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -221,7 +221,18 @@ import { } from '../storage/git.js'; import { isGitNexusManagedPath } from '../storage/gitnexus-managed-paths.js'; import { getMaxFileSizeBytes } from './ingestion/utils/max-file-size.js'; -import type { CachedEmbedding } from './embeddings/types.js'; +import { + cacheRowCount, + discardScopedEmbeddingSpills, + disposeEmbeddingSpill, + withEmbeddingSpillScope, + emptyCachedEmbeddingsSnapshot, + EmbeddingSpillReader, + materializeCachedEmbeddings, + normalizeCachedEmbeddings, + snapshotEmbeddingDims, + type CachedEmbeddingsSnapshot, +} from './embeddings/embedding-restore-spill.js'; import { generateAIContextFiles } from '../cli/ai-context.js'; import { sanitizeDetectedBranch } from '../cli/analyze-config.js'; import { @@ -1166,59 +1177,64 @@ export async function runFullAnalysis( let writeTarget = await resolveWriteTarget(repoPath, options); let lock = await acquireIndexLock(writeTarget.metaDir, acquireOpts); - try { - requireExclusiveIndexLock( - lock, - `Cannot acquire the index lock at ${writeTarget.metaDir}; refusing an unlocked analysis.`, - ); - // #2658 review H2: acquireIndexLock can wait up to the timeout ceiling, - // during which git HEAD/branch — and thus the resolved write slot — may - // change (a commit lands, a branch is switched, or another writer adopts the - // flat slot). The pre-wait snapshot must NOT be reused: re-resolve UNDER the - // lock so the freshness check (`existingMeta.lastCommit === currentCommit`) - // and the meta stamps see current git state, honoring the module's "re-check - // freshness after acquiring" contract. If the slot itself moved we hold the - // WRONG lock — release and re-acquire the correct one. Bounded so a - // pathologically churning checkout can't loop forever; after the cap we - // proceed on the current lock. The loop is INSIDE the try so a re-resolve - // that throws (e.g. a `--branch` that stopped matching the now-switched - // checkout) still releases the held lock via `finally` (no leak). - const MAX_RELOCK = 3; - for (let attempt = 0; attempt < MAX_RELOCK; attempt++) { - // Never pass the pre-lock storagePath as already-validated: requireStoragePath - // must run again under the lock so a now-foreign slot aborts (and finally - // still releases the lock). - const fresh = await resolveWriteTarget(repoPath, options); - if (fresh.metaDir === writeTarget.metaDir) { - writeTarget = fresh; // same slot — adopt the freshly-read commit/branch/placement - break; - } - log( - `Index write target moved while waiting for the lock ` + - `(${writeTarget.metaDir} → ${fresh.metaDir}); re-acquiring the correct slot.`, - ); - lock.release(); - writeTarget = fresh; - lock = await acquireIndexLock(fresh.metaDir, acquireOpts); + return withEmbeddingSpillScope(async () => { + try { requireExclusiveIndexLock( lock, - `Cannot acquire the index lock at ${fresh.metaDir}; refusing an unlocked analysis.`, + `Cannot acquire the index lock at ${writeTarget.metaDir}; refusing an unlocked analysis.`, ); - if (attempt === MAX_RELOCK - 1) { - log('Index write target still moving after repeated re-acquire; proceeding on this lock.'); + // #2658 review H2: acquireIndexLock can wait up to the timeout ceiling, + // during which git HEAD/branch — and thus the resolved write slot — may + // change (a commit lands, a branch is switched, or another writer adopts the + // flat slot). The pre-wait snapshot must NOT be reused: re-resolve UNDER the + // lock so the freshness check (`existingMeta.lastCommit === currentCommit`) + // and the meta stamps see current git state, honoring the module's "re-check + // freshness after acquiring" contract. If the slot itself moved we hold the + // WRONG lock — release and re-acquire the correct one. Bounded so a + // pathologically churning checkout can't loop forever; after the cap we + // proceed on the current lock. The loop is INSIDE the try so a re-resolve + // that throws (e.g. a `--branch` that stopped matching the now-switched + // checkout) still releases the held lock via `finally` (no leak). + const MAX_RELOCK = 3; + for (let attempt = 0; attempt < MAX_RELOCK; attempt++) { + // Never pass the pre-lock storagePath as already-validated: requireStoragePath + // must run again under the lock so a now-foreign slot aborts (and finally + // still releases the lock). + const fresh = await resolveWriteTarget(repoPath, options); + if (fresh.metaDir === writeTarget.metaDir) { + writeTarget = fresh; // same slot — adopt the freshly-read commit/branch/placement + break; + } + log( + `Index write target moved while waiting for the lock ` + + `(${writeTarget.metaDir} → ${fresh.metaDir}); re-acquiring the correct slot.`, + ); + lock.release(); + writeTarget = fresh; + lock = await acquireIndexLock(fresh.metaDir, acquireOpts); + requireExclusiveIndexLock( + lock, + `Cannot acquire the index lock at ${fresh.metaDir}; refusing an unlocked analysis.`, + ); + if (attempt === MAX_RELOCK - 1) { + log( + 'Index write target still moving after repeated re-acquire; proceeding on this lock.', + ); + } } + return await runFullAnalysisInner( + repoPath, + options, + callbacks, + writeTarget, + contentRetention, + runnerIdentityAtBootstrap, + ); + } finally { + discardScopedEmbeddingSpills(); + lock.release(); } - return await runFullAnalysisInner( - repoPath, - options, - callbacks, - writeTarget, - contentRetention, - runnerIdentityAtBootstrap, - ); - } finally { - lock.release(); - } + }); } async function runFullAnalysisInner( @@ -2201,8 +2217,18 @@ async function runFullAnalysisInner( // The default-preserve branch is what makes a routine `analyze` (e.g. a // post-commit hook) safe: a multi-minute embedding pass is no longer // silently dropped just because the caller omitted `--embeddings`. - let cachedEmbeddingNodeIds = new Set(); - let cachedEmbeddings: CachedEmbedding[] = []; + let cachedSnapshot: CachedEmbeddingsSnapshot = emptyCachedEmbeddingsSnapshot(); + const adoptCachedEmbeddings = (raw: CachedEmbeddingsSnapshot): void => { + cachedSnapshot = normalizeCachedEmbeddings(raw); + }; + const discardCachedEmbeddings = (): void => { + disposeEmbeddingSpill(cachedSnapshot.spill); + cachedSnapshot = emptyCachedEmbeddingsSnapshot(); + }; + const discardCachedEmbeddingSpill = (): void => { + disposeEmbeddingSpill(cachedSnapshot.spill); + cachedSnapshot = { ...cachedSnapshot, spill: undefined }; + }; const existingEmbeddingCount = existingMeta?.stats?.embeddings ?? 0; const { @@ -2237,7 +2263,7 @@ async function runFullAnalysisInner( // of the predicted `willTryIncremental`). The post-pipeline branch may // disagree with the prediction (e.g. when the pipeline produces zero // File nodes, `isIncremental` flips false and the full-rebuild path - // wipes the DB) — loading unconditionally is cheap insurance against + // wipes the DB) — loading unconditionally is insurance against // silently dropping embeddings on a mispredicted run. The re-insert // step gates itself on the actual `isIncremental` value to avoid // PK-conflicts when the incremental writeback path keeps the rows. @@ -2251,9 +2277,7 @@ async function runFullAnalysisInner( try { progress('embeddings', 0, 'Caching embeddings...'); await initAnalysisLbug(lbugPath); - const cached = await loadCachedEmbeddings(); - cachedEmbeddingNodeIds = cached.embeddingNodeIds; - cachedEmbeddings = cached.embeddings; + adoptCachedEmbeddings(await loadCachedEmbeddings()); await closeLbug(); } catch (err: any) { // Surface cache-load failures explicitly: silently swallowing here would @@ -2264,8 +2288,7 @@ async function runFullAnalysisInner( `(${err?.message ?? String(err)}). ` + `Embeddings will not be preserved on this run.`, ); - cachedEmbeddingNodeIds = new Set(); - cachedEmbeddings = []; + discardCachedEmbeddings(); try { await closeLbug(); } catch { @@ -2378,6 +2401,7 @@ async function runFullAnalysisInner( }, ); } catch (err) { + discardCachedEmbeddingSpill(); await removeColdParseRebuildDir(coldParseRebuildDir, true); throw err; } @@ -2614,6 +2638,7 @@ async function runFullAnalysisInner( try { await wipeLbugDbFiles(buildPath); } catch (error) { + discardCachedEmbeddingSpill(); if (liveIndexMutationStarted) recordLiveIndexMutationRisk(error); throw error; } @@ -2642,6 +2667,7 @@ async function runFullAnalysisInner( try { await initAnalysisLbug(buildPath); } catch (error) { + discardCachedEmbeddingSpill(); if (liveIndexMutationStarted) recordLiveIndexMutationRisk(error); throw error; } @@ -2963,7 +2989,7 @@ async function runFullAnalysisInner( const extensionForcedRebuild = !embeddingRowDmlSafe || !ftsRowDmlSafe; // `!options.dropEmbeddings` (H1): this rescue reads the rows back OUT of // the DB, so it must never fire on the one path whose entire purpose is to - // destroy them. `--drop-embeddings` deliberately leaves `cachedEmbeddings` + // destroy them. `--drop-embeddings` deliberately leaves `cachedSnapshot` // empty (`deriveEmbeddingMode` returns `shouldLoadCache: false` for it by // construction — see the four-mode comment at the cache-load site), and its // `options.force = true` conversion sits INSIDE @@ -2978,12 +3004,16 @@ async function runFullAnalysisInner( // while rows survive ⇒ `hasExisting` false ⇒ `shouldLoadCache` false), i.e. // it would fix the wipe by deleting the safeguard. Covers // `--drop-embeddings --embeddings` too — the rescue repopulates - // `cachedEmbeddingNodeIds`, which Phase 4 hands `runEmbeddingPipeline` as + // `cachedSnapshot.embeddingNodeIds`, which Phase 4 hands `runEmbeddingPipeline` as // the already-embedded set, so the very nodes the user asked to REGENERATE // would be skipped. - if (extensionForcedRebuild && !options.dropEmbeddings && cachedEmbeddings.length === 0) { + if ( + extensionForcedRebuild && + !options.dropEmbeddings && + cacheRowCount(cachedSnapshot) === 0 + ) { // The escalation below WIPES the DB files, and Phase 3.5 restores - // embedding rows from `cachedEmbeddings` — which is only populated when + // embedding rows from `cachedSnapshot` — which is only populated when // `deriveEmbeddingMode` saw `meta.stats.embeddings > 0`. A DB whose meta // under-reports its embeddings (meta restored from an older run, or a // count that never got stamped) would therefore have every vector @@ -2991,14 +3021,21 @@ async function runFullAnalysisInner( // while the DB is still intact — a plain MATCH, which needs no VECTOR // extension. Rows whose owning node is gone are dropped by Phase 3.5's // live-graph filter, exactly as on any other wiped path. - const rescued = await loadCachedEmbeddings(); - if (rescued.embeddings.length > 0) { - cachedEmbeddings = rescued.embeddings; - cachedEmbeddingNodeIds = rescued.embeddingNodeIds; + try { + adoptCachedEmbeddings(await loadCachedEmbeddings()); + if (cacheRowCount(cachedSnapshot) > 0) { + log( + `Preserving ${cacheRowCount(cachedSnapshot)} embedding row(s) across the forced rebuild ` + + `(the index metadata did not account for them).`, + ); + } + } catch (err: any) { log( - `Preserving ${rescued.embeddings.length} embedding row(s) across the forced rebuild ` + - `(the index metadata did not account for them).`, + `Warning: could not load cached embeddings ` + + `(${err?.message ?? String(err)}). ` + + `Embeddings will not be preserved on this run.`, ); + discardCachedEmbeddings(); } } // Hoisted out of the `||` below (§5.D): the size verdict has to be KNOWN @@ -3624,26 +3661,27 @@ async function runFullAnalysisInner( // propagates errors (a completed writeback means a deterministic // delete outcome) and this process holds the exclusive DB lock (no // concurrent writer). - // The per-batch try/catch stays as a last-resort guard only — it no - // longer fires on the happy path. + // Materialize runs outside the insert catch so a spill I/O failure is not + // treated as a benign PK conflict. Any node with a failed restore batch is + // marked stale in the Phase 4 map so leftover chunks are deleted and rembedded. let restoredEmbeddingCount = 0; - if (cachedEmbeddings.length > 0) { - const cachedDims = cachedEmbeddings[0].embedding.length; + const restoreFailedNodeIds = new Set(); + if (cacheRowCount(cachedSnapshot) > 0) { + const cachedDims = snapshotEmbeddingDims(cachedSnapshot); const { EMBEDDING_DIMS } = await import('./lbug/schema.js'); - if (cachedDims !== EMBEDDING_DIMS) { + if (cachedDims !== undefined && cachedDims !== EMBEDDING_DIMS) { // Dimensions changed (e.g. switched embedding model) — discard cache and re-embed all log( `Embedding dimensions changed (${cachedDims}d -> ${EMBEDDING_DIMS}d), discarding cache`, ); - cachedEmbeddings = []; - cachedEmbeddingNodeIds = new Set(); + discardCachedEmbeddings(); } else { const { batchInsertEmbeddings: batchInsert } = await import('./embeddings/embedding-pipeline.js'); // (1) Live-graph filter — the FULL pipeline graph (always produced), // NOT the incremental subgraph, or unchanged files' rows would be // dropped from the restore set. - const liveEmbeddings = cachedEmbeddings.filter( + const liveEmbeddings = cachedSnapshot.rows.filter( (e) => pipelineResult.graph.getNode(e.nodeId) !== undefined, ); // (2) Restore-scope filter (see the discipline note above). @@ -3656,13 +3694,36 @@ async function runFullAnalysisInner( }); progress('embeddings', 88, `Restoring ${rowsToRestore.length} cached embeddings...`); const EMBED_BATCH = 200; - for (const batch of chunk(rowsToRestore, EMBED_BATCH)) { - try { - await batchInsert(executeWithReusedStatement, batch); - restoredEmbeddingCount += batch.length; - } catch { - /* last-resort guard — conflict-free by construction above */ + let spillReader: EmbeddingSpillReader | undefined; + try { + for (const batch of chunk(rowsToRestore, EMBED_BATCH)) { + let materialized; + try { + if (!spillReader && cachedSnapshot.spill && cachedSnapshot.embeddings.length === 0) { + spillReader = new EmbeddingSpillReader(cachedSnapshot.spill); + } + materialized = materializeCachedEmbeddings(cachedSnapshot, batch, spillReader); + } catch (err) { + for (const row of batch) restoreFailedNodeIds.add(row.nodeId); + log( + `Warning: could not materialize ${batch.length} cached embedding(s) for restore ` + + `(${(err as Error).message}); those nodes will be re-embedded if this run generates embeddings.`, + ); + continue; + } + try { + await batchInsert(executeWithReusedStatement, materialized); + restoredEmbeddingCount += batch.length; + } catch (err) { + for (const row of batch) restoreFailedNodeIds.add(row.nodeId); + log( + `Warning: could not restore ${batch.length} cached embedding(s) ` + + `(${(err as Error).message}); those nodes will be re-embedded if this run generates embeddings.`, + ); + } } + } finally { + spillReader?.close(); } // Legacy-orphan sweep (FIX 3, finder B): the live-graph filter's @@ -3679,7 +3740,7 @@ async function runFullAnalysisInner( // sweep failure must never fail a completed writeback, so the whole // sweep warns-and-continues. if (deletedFilePathsForRestore !== null) { - const orphanRowIds = cachedEmbeddings + const orphanRowIds = cachedSnapshot.rows .filter((e) => pipelineResult.graph.getNode(e.nodeId) === undefined) .map((e) => `${e.nodeId}:${e.chunkIndex}`); if (orphanRowIds.length > 0) { @@ -3708,6 +3769,9 @@ async function runFullAnalysisInner( } } } + // Vectors are on disk only to survive the wipe/delete. After restore, + // drop the spill so Phase 4 does not keep a multi-GB temp file open. + discardCachedEmbeddingSpill(); // ── Phase 4: Embeddings (90–98%) ────────────────────────────────── const stats = await getLbugStats(); @@ -3957,9 +4021,16 @@ async function runFullAnalysisInner( const embeddingIdentity = embeddingIdentityForRun; // Build a Map from cached embeddings for incremental mode let existingEmbeddings: Map | undefined; - if (cachedEmbeddingNodeIds.size > 0) { + if (cachedSnapshot.embeddingNodeIds.size > 0) { existingEmbeddings = new Map(); - for (const e of cachedEmbeddings) { + for (const e of cachedSnapshot.rows) { + if (restoreFailedNodeIds.has(e.nodeId)) { + // Any failed batch for this node: mark stale so Phase 4 DELETEs + // leftover chunks and re-embeds. Omitting the id would treat the + // node as new and PK-conflict on rows that already restored. + existingEmbeddings.set(e.nodeId, STALE_HASH_SENTINEL); + continue; + } existingEmbeddings.set(e.nodeId, e.contentHash ?? STALE_HASH_SENTINEL); } } @@ -4036,7 +4107,7 @@ async function runFullAnalysisInner( progress('embeddings', scaled, label); }, {}, - cachedEmbeddingNodeIds.size > 0 ? cachedEmbeddingNodeIds : undefined, + cachedSnapshot.embeddingNodeIds.size > 0 ? cachedSnapshot.embeddingNodeIds : undefined, existingEmbeddings, { forceReembedNodeIds: pendingEmbeddingNodeIds, @@ -4718,6 +4789,7 @@ async function runFullAnalysisInner( } } await removeColdParseRebuildDir(coldParseRebuildDir, true); + discardCachedEmbeddingSpill(); if (liveIndexMutationStarted) { // Preserve the original error identity/prototype: callers distinguish // IndexLockTimeoutError and other domain failures with `instanceof`. diff --git a/gitnexus/test/integration/load-cached-embeddings-spill.test.ts b/gitnexus/test/integration/load-cached-embeddings-spill.test.ts new file mode 100644 index 000000000..99adb47d0 --- /dev/null +++ b/gitnexus/test/integration/load-cached-embeddings-spill.test.ts @@ -0,0 +1,113 @@ +/** + * Real-DB coverage for #3306: loadCachedEmbeddings must stream CodeEmbedding + * rows instead of getAll()+map(Number) of the whole table, and must be able + * to spill vectors so incremental analyze does not keep every embedding in + * the V8 heap. + */ +import fs from 'node:fs'; +import path from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { createTempDir, type TestDBHandle } from '../helpers/test-db.js'; +import { EMBEDDING_DIMS } from '../../src/core/lbug/schema.js'; +import { batchInsertEmbeddings } from '../../src/core/embeddings/embedding-pipeline.js'; +import { + disposeEmbeddingSpill, + materializeCachedEmbeddings, +} from '../../src/core/embeddings/embedding-restore-spill.js'; + +describe('loadCachedEmbeddings streaming (#3306)', () => { + let tmp: TestDBHandle | undefined; + + afterEach(async () => { + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + try { + await adapter.closeLbug(); + } catch { + /* already closed */ + } + await tmp?.cleanup(); + tmp = undefined; + }); + + async function seedDb(rowCount: number) { + tmp = await createTempDir('gitnexus-lbug-'); + const dbPath = path.join(tmp.dbPath, 'lbug'); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + await adapter.initLbug(dbPath); + const rows = Array.from({ length: rowCount }, (_, i) => ({ + nodeId: `Function:src/f${i}.ts:fn${i}:1`, + chunkIndex: 0, + startLine: 1, + endLine: 3, + embedding: Array.from({ length: EMBEDDING_DIMS }, (__, d) => (d === 0 ? i + 1 : 0)), + contentHash: `hash-${i}`, + })); + await batchInsertEmbeddings(adapter.executeWithReusedStatement, rows); + return { adapter, rows }; + } + + it('materializes a small table in RAM (skip-fts / mock-compatible shape)', async () => { + const { adapter, rows } = await seedDb(3); + const cached = await adapter.loadCachedEmbeddings(); + expect(cached.spill).toBeUndefined(); + expect(cached.embeddings).toHaveLength(3); + expect(cached.rows).toHaveLength(3); + expect(cached.embeddingNodeIds.size).toBe(3); + expect(cached.embeddings.map((e) => e.nodeId).sort()).toEqual(rows.map((r) => r.nodeId).sort()); + expect(cached.embeddings.find((e) => e.nodeId === rows[1]!.nodeId)?.embedding[0]).toBe(2); + }); + + it('streams into a spill file when the in-memory limit is 0 and restores a subset', async () => { + const { adapter, rows } = await seedDb(12); + const cached = await adapter.loadCachedEmbeddings({ inMemoryRowLimit: 0 }); + try { + expect(cached.embeddings).toEqual([]); + expect(cached.spill?.rowCount).toBe(12); + expect(cached.rows).toHaveLength(12); + const wanted = new Set([rows[0]!.nodeId, rows[5]!.nodeId, rows[10]!.nodeId]); + const subset = materializeCachedEmbeddings( + cached, + cached.rows.filter((meta) => wanted.has(meta.nodeId)), + ); + expect(subset).toHaveLength(3); + const byId = new Map(subset.map((row) => [row.nodeId, row])); + expect(byId.get(rows[0]!.nodeId)?.embedding[0]).toBe(1); + expect(byId.get(rows[5]!.nodeId)?.embedding[0]).toBe(6); + expect(byId.get(rows[10]!.nodeId)?.embedding[0]).toBe(11); + expect(byId.get(rows[0]!.nodeId)?.contentHash).toBe(rows[0]!.contentHash); + } finally { + disposeEmbeddingSpill(cached.spill); + } + }); + + it('flips from RAM to spill once a non-zero in-memory limit is crossed', async () => { + const { adapter, rows } = await seedDb(8); + const cached = await adapter.loadCachedEmbeddings({ inMemoryRowLimit: 4 }); + try { + expect(cached.embeddings).toEqual([]); + expect(cached.spill?.rowCount).toBe(8); + expect(cached.rows).toHaveLength(8); + expect(fs.statSync(cached.spill!.path).size).toBe(12 + 8 * EMBEDDING_DIMS * 4); + const wanted = new Set([rows[2]!.nodeId, rows[3]!.nodeId]); + const subset = materializeCachedEmbeddings( + cached, + cached.rows.filter((meta) => wanted.has(meta.nodeId)), + ); + expect(subset).toHaveLength(2); + const byId = new Map(subset.map((row) => [row.nodeId, row])); + expect(byId.get(rows[2]!.nodeId)?.embedding[0]).toBe(3); + expect(byId.get(rows[3]!.nodeId)?.embedding[0]).toBe(4); + } finally { + disposeEmbeddingSpill(cached.spill); + } + }); + + it('surfaces a spill write failure instead of adopting an empty snapshot', async () => { + const { adapter } = await seedDb(3); + const spillDir = path.join(tmp!.dbPath, 'not-a-directory'); + fs.writeFileSync(spillDir, 'x'); + await expect(adapter.loadCachedEmbeddings({ inMemoryRowLimit: 0, spillDir })).rejects.toThrow( + /ENOTDIR|not a directory|ENOSPC|EACCES/i, + ); + }); +}); diff --git a/gitnexus/test/unit/embedding-restore-spill.test.ts b/gitnexus/test/unit/embedding-restore-spill.test.ts new file mode 100644 index 000000000..b11695110 --- /dev/null +++ b/gitnexus/test/unit/embedding-restore-spill.test.ts @@ -0,0 +1,274 @@ +import { existsSync, statSync, writeFileSync } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { afterEach, describe, expect, it, vi } from 'vitest'; +import { + abortCachedEmbeddingsBuilder, + cacheRowCount, + createCachedEmbeddingsBuilder, + DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT, + discardLiveEmbeddingSpills, + discardScopedEmbeddingSpills, + disposeEmbeddingSpill, + EmbeddingSpillReader, + finalizeCachedEmbeddingsSnapshot, + ingestCachedEmbeddingRow, + materializeCachedEmbeddings, + normalizeCachedEmbeddings, + readSpillVectors, + resolveEmbeddingCacheInMemoryRowLimit, + snapshotEmbeddingDims, + withEmbeddingSpillScope, +} from '../../src/core/embeddings/embedding-restore-spill.js'; + +const DIMS = 8; + +function vector(fill: number): number[] { + return Array.from({ length: DIMS }, () => fill); +} + +function row(id: string, fill: number, hash = `hash-${id}`) { + return { + nodeId: id, + chunkIndex: 0, + startLine: 1, + endLine: 2, + embedding: vector(fill), + contentHash: hash, + }; +} + +describe('embedding-restore-spill (#3306)', () => { + const spills: Array<{ path: string }> = []; + afterEach(() => { + for (const spill of spills) disposeEmbeddingSpill(spill); + spills.length = 0; + }); + + it('keeps small tables in RAM and does not leave a spill file', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 4, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('n1', 0.25), true); + ingestCachedEmbeddingRow(builder, row('n2', 0.5), true); + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + expect(snapshot.spill).toBeUndefined(); + expect(existsSync(builder.writer.path)).toBe(false); + expect(snapshot.embeddings).toHaveLength(2); + expect(snapshot.rows).toHaveLength(2); + expect(snapshot.embeddings[0]?.embedding[0]).toBeCloseTo(0.25); + expect(cacheRowCount(snapshot)).toBe(2); + expect(snapshotEmbeddingDims(snapshot)).toBe(DIMS); + }); + + it('spills vectors once the in-memory limit is exceeded and materializes a subset', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 2, + spillDir: os.tmpdir(), + }); + for (let i = 0; i < 5; i++) { + ingestCachedEmbeddingRow(builder, row(`n${i}`, i + 1), true); + } + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + if (snapshot.spill) spills.push(snapshot.spill); + expect(snapshot.embeddings).toEqual([]); + expect(snapshot.spill?.rowCount).toBe(5); + expect(existsSync(snapshot.spill!.path)).toBe(true); + expect(snapshot.embeddingNodeIds.size).toBe(5); + + const subset = materializeCachedEmbeddings(snapshot, snapshot.rows.slice(1, 3)); + expect(subset).toHaveLength(2); + expect(subset[0]?.nodeId).toBe('n1'); + expect(subset[0]?.embedding[0]).toBeCloseTo(2); + expect(subset[1]?.embedding[0]).toBeCloseTo(3); + }); + + it('always spills when the in-memory limit is 0 (no Number[] table in RAM)', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('only', 0.75), true); + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + if (snapshot.spill) spills.push(snapshot.spill); + expect(snapshot.embeddings).toEqual([]); + expect(snapshot.rows).toHaveLength(1); + expect(materializeCachedEmbeddings(snapshot, snapshot.rows)[0]?.embedding[0]).toBeCloseTo(0.75); + }); + + it('aborts an unfinished builder without leaking a spill file', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('n1', 1), true); + expect(existsSync(builder.writer.path)).toBe(true); + abortCachedEmbeddingsBuilder(builder); + expect(existsSync(builder.writer.path)).toBe(false); + }); + + it('normalizes mock {embeddings} payloads so Phase 3.5 can restore without a spill', () => { + const snapshot = normalizeCachedEmbeddings({ + embeddingNodeIds: new Set(['Function:a:foo']), + embeddings: [ + { + nodeId: 'Function:a:foo', + chunkIndex: 0, + startLine: 0, + endLine: 3, + embedding: vector(0.1), + contentHash: 'stub', + }, + ], + }); + expect(snapshot.rows).toHaveLength(1); + expect(snapshot.rows[0]?.vectorIndex).toBe(0); + const restored = materializeCachedEmbeddings(snapshot, snapshot.rows); + expect(restored[0]?.contentHash).toBe('stub'); + expect(restored[0]?.embedding).toHaveLength(DIMS); + }); + + it('defaults the in-memory row limit to 2048 and honors GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', () => { + expect(DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT).toBe(2048); + vi.stubEnv('GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', ''); + expect(resolveEmbeddingCacheInMemoryRowLimit()).toBe(2048); + vi.stubEnv('GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', '0'); + expect(resolveEmbeddingCacheInMemoryRowLimit()).toBe(0); + vi.stubEnv('GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', '12'); + expect(resolveEmbeddingCacheInMemoryRowLimit()).toBe(12); + vi.stubEnv('GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', 'nope'); + expect(resolveEmbeddingCacheInMemoryRowLimit()).toBe(2048); + vi.unstubAllEnvs(); + expect(resolveEmbeddingCacheInMemoryRowLimit(7)).toBe(7); + expect(resolveEmbeddingCacheInMemoryRowLimit(-1)).toBe(2048); + }); + + it('spills above the default 2048-row limit with header-plus-body size 12 + N * D * 4', () => { + const n = DEFAULT_EMBEDDING_CACHE_IN_MEMORY_ROW_LIMIT + 1; + const builder = createCachedEmbeddingsBuilder({ spillDir: os.tmpdir() }); + for (let i = 0; i < n; i++) { + ingestCachedEmbeddingRow(builder, row(`n${i}`, 1), true); + } + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + if (snapshot.spill) spills.push(snapshot.spill); + expect(snapshot.embeddings).toEqual([]); + expect(snapshot.spill?.rowCount).toBe(n); + expect(statSync(snapshot.spill!.path).size).toBe(12 + n * DIMS * 4); + }); + + it('rejects a dim mismatch once spilling and rejects a short or bad-magic header', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('a', 1), true); + expect(() => + ingestCachedEmbeddingRow(builder, { ...row('b', 2), embedding: [1, 2, 3] }, true), + ).toThrow(/dim mismatch/); + abortCachedEmbeddingsBuilder(builder); + + const shortPath = path.join(os.tmpdir(), `gitnexus-embed-restore-short-${process.pid}.bin`); + writeFileSync(shortPath, Buffer.from('NOPE')); + spills.push({ path: shortPath }); + expect(() => readSpillVectors({ path: shortPath, dims: DIMS, rowCount: 1 }, [0])).toThrow( + /invalid embedding spill header/, + ); + + const badMagic = Buffer.alloc(12); + badMagic.write('NOPE', 0, 4, 'ascii'); + badMagic.writeUInt8(1, 4); + badMagic.writeUInt32LE(DIMS, 5); + const badPath = path.join(os.tmpdir(), `gitnexus-embed-restore-bad-${process.pid}.bin`); + writeFileSync(badPath, badMagic); + spills.push({ path: badPath }); + expect(() => readSpillVectors({ path: badPath, dims: DIMS, rowCount: 1 }, [0])).toThrow( + /invalid embedding spill header/, + ); + }); + + it('reuses an open spill reader across materialize batches', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + for (let i = 0; i < 4; i++) { + ingestCachedEmbeddingRow(builder, row(`n${i}`, i + 1), true); + } + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + if (snapshot.spill) spills.push(snapshot.spill); + const reader = new EmbeddingSpillReader(snapshot.spill!); + try { + const first = materializeCachedEmbeddings(snapshot, snapshot.rows.slice(0, 2), reader); + const second = materializeCachedEmbeddings(snapshot, snapshot.rows.slice(2, 4), reader); + expect(first[0]?.embedding[0]).toBeCloseTo(1); + expect(second[1]?.embedding[0]).toBeCloseTo(4); + } finally { + reader.close(); + } + }); + + it('scoped discard unlinks only spills created in that analyze run', async () => { + const other = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(other, row('other', 1), true); + const otherSnapshot = finalizeCachedEmbeddingsSnapshot(other); + if (otherSnapshot.spill) spills.push(otherSnapshot.spill); + expect(existsSync(otherSnapshot.spill!.path)).toBe(true); + + await withEmbeddingSpillScope(async () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('scoped', 2), true); + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + expect(existsSync(snapshot.spill!.path)).toBe(true); + discardScopedEmbeddingSpills(); + expect(existsSync(snapshot.spill!.path)).toBe(false); + expect(existsSync(otherSnapshot.spill!.path)).toBe(true); + }); + }); + + it('unlinks a finished spill that was not disposed', () => { + const builder = createCachedEmbeddingsBuilder({ + inMemoryRowLimit: 0, + spillDir: os.tmpdir(), + }); + ingestCachedEmbeddingRow(builder, row('n1', 1), true); + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + expect(snapshot.spill).toBeDefined(); + expect(existsSync(snapshot.spill!.path)).toBe(true); + discardLiveEmbeddingSpills(); + expect(existsSync(snapshot.spill!.path)).toBe(false); + }); + + it('throws when materializing a meta row with no matching vector', () => { + const snapshot = normalizeCachedEmbeddings({ + embeddings: [ + { + nodeId: 'Function:a:foo', + chunkIndex: 0, + startLine: 0, + endLine: 3, + embedding: vector(0.1), + contentHash: 'stub', + }, + ], + }); + expect(() => + materializeCachedEmbeddings(snapshot, [ + { + nodeId: 'Function:missing:bar', + chunkIndex: 0, + startLine: 0, + endLine: 1, + contentHash: 'x', + vectorIndex: 99, + }, + ]), + ).toThrow(/missing cached embedding Function:missing:bar:0/); + }); +}); diff --git a/gitnexus/test/unit/incremental-dirty-recovery.test.ts b/gitnexus/test/unit/incremental-dirty-recovery.test.ts index eafb933ce..0192cac1c 100644 --- a/gitnexus/test/unit/incremental-dirty-recovery.test.ts +++ b/gitnexus/test/unit/incremental-dirty-recovery.test.ts @@ -19,7 +19,7 @@ */ import { writeFile, readFile } from 'fs/promises'; -import { describe, it, expect } from 'vitest'; +import { afterEach, describe, expect, it, vi } from 'vitest'; import { getStoragePaths, saveMeta, @@ -35,92 +35,109 @@ import { seedEmbeddingsForFiles } from '../helpers/embedding-seed.js'; const setupMiniRepo = () => setupSharedMiniRepo('gitnexus-incr-dirty-rec-'); +async function parkCrashSidecarsAndRebuild(options?: { alwaysSpill?: boolean }) { + vi.stubEnv('GITNEXUS_WORKER_READY_TIMEOUT_MS', '60000'); + if (options?.alwaysSpill) { + vi.stubEnv('GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT', '0'); + } + const repo = await setupMiniRepo(); + try { + const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); + await runFullAnalysis(repo.dbPath, { skipAgentsMd: true }, { onProgress: () => {} }); + + // Seed real embeddings BEFORE the tamper (tri-review 4669518496 / U5): + // with meta.stats.embeddings = 0 the recovery run derived + // shouldLoadCache=false and never opened the DB pre-wipe — this test + // was vacuous about the exact open the parking protects. Seeded rows + + // a stats stamp route the recovery (which runs force:true internally, + // so forceRegenerate → shouldLoadCache) through the REAL + // embedding-cache preservation open on the just-parked DB. + const { storagePath, lbugPath } = getStoragePaths(repo.dbPath); + const seededIdsByFile = await seedEmbeddingsForFiles( + repo.dbPath, + ['src/handler.ts', 'src/logger.ts'], + 1, + ); + const seededNodeIds = [...seededIdsByFile.values()].flat(); + expect(seededNodeIds.length).toBeGreaterThan(0); + + // Simulate a crashed incremental writeback: dirty flag in meta plus + // leftover sidecars whose bytes must never be replayed. 8KB puts the + // WAL above the tiny-orphan threshold — the state the sidecar + // preflight deliberately leaves in place for engine replay. + const meta = await loadMeta(storagePath); + const tampered: RepoMeta = { + ...meta!, + stats: { ...meta!.stats, embeddings: seededNodeIds.length }, + incrementalInProgress: { + startedAt: Date.now() - 60_000, + toWriteCount: 12, + phase: 'load-graph', + }, + }; + await saveMeta(storagePath, tampered); + const walGarbage = Buffer.alloc(8192, 0xab); + const shadowGarbage = Buffer.alloc(4096, 0xcd); + await writeFile(`${lbugPath}.wal`, walGarbage); + await writeFile(`${lbugPath}.shadow`, shadowGarbage); + + const logs: string[] = []; + // embeddingsNodeLimit: 1 (KTD9): the recovery runs force:true + // internally, and the seeded stats would otherwise route Phase 4 into + // a real embedder in CI — the 1-node cap suppresses generation while + // leaving the preserve/restore path fully live. On linux the + // wipe-and-restore vector-index seam then fires for real (statically + // linked VECTOR): a CREATE_VECTOR_INDEX over the restored rows is + // expected and harmless here. + const recovered = await runFullAnalysis( + repo.dbPath, + { skipAgentsMd: true, embeddingsNodeLimit: 1 }, + { onProgress: () => {}, onLog: (m) => logs.push(m) }, + ); + expect(recovered.alreadyUpToDate).toBeUndefined(); + + // Both sidecars were parked verbatim (renamed, never deleted) before + // any open could replay them… + expect(Buffer.compare(await readFile(`${lbugPath}.wal.dirty-recovery`), walGarbage)).toBe(0); + expect(Buffer.compare(await readFile(`${lbugPath}.shadow.dirty-recovery`), shadowGarbage)).toBe( + 0, + ); + const joinedLogs = logs.join('\n'); + expect(joinedLogs).toContain('Parked lbug.wal.dirty-recovery, lbug.shadow.dirty-recovery'); + + // …the run traversed the REAL pre-wipe preservation open — recovery's + // internal force on an embedded repo upgrades to regenerate mode, whose + // banner only prints when existingEmbeddingCount was read from the + // seeded stats and the cache-load path engaged… + expect(joinedLogs).toContain( + `--force on a repo with ${seededNodeIds.length} existing embeddings`, + ); + // …with generation itself cap-suppressed (no embedder in CI): + expect(joinedLogs).toContain('exceeds the 1-node safety cap'); + + // …and the rebuild completed into a clean index: dirty flag cleared, + // and the seeded embeddings survived the park → open → wipe → restore + // round-trip (the strongest signal the preservation open really ran: + // the DB was wiped, so these rows can only come from the cache load). + const after = await loadMeta(storagePath); + expect(after!.incrementalInProgress).toBeUndefined(); + expect(after!.stats?.embeddings).toBe(seededNodeIds.length); + } finally { + vi.unstubAllEnvs(); + await repo.cleanup(); + } +} + describe('runFullAnalysis — dirty-flag recovery sidecar parking (#2409)', () => { + afterEach(() => { + vi.unstubAllEnvs(); + }); + it('parks the crashed run WAL/shadow sidecars before reopening, then rebuilds clean', async () => { - const repo = await setupMiniRepo(); - try { - const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); - await runFullAnalysis(repo.dbPath, { skipAgentsMd: true }, { onProgress: () => {} }); + await parkCrashSidecarsAndRebuild(); + }, 300_000); - // Seed real embeddings BEFORE the tamper (tri-review 4669518496 / U5): - // with meta.stats.embeddings = 0 the recovery run derived - // shouldLoadCache=false and never opened the DB pre-wipe — this test - // was vacuous about the exact open the parking protects. Seeded rows + - // a stats stamp route the recovery (which runs force:true internally, - // so forceRegenerate → shouldLoadCache) through the REAL - // embedding-cache preservation open on the just-parked DB. - const { storagePath, lbugPath } = getStoragePaths(repo.dbPath); - const seededIdsByFile = await seedEmbeddingsForFiles( - repo.dbPath, - ['src/handler.ts', 'src/logger.ts'], - 1, - ); - const seededNodeIds = [...seededIdsByFile.values()].flat(); - expect(seededNodeIds.length).toBeGreaterThan(0); - - // Simulate a crashed incremental writeback: dirty flag in meta plus - // leftover sidecars whose bytes must never be replayed. 8KB puts the - // WAL above the tiny-orphan threshold — the state the sidecar - // preflight deliberately leaves in place for engine replay. - const meta = await loadMeta(storagePath); - const tampered: RepoMeta = { - ...meta!, - stats: { ...meta!.stats, embeddings: seededNodeIds.length }, - incrementalInProgress: { - startedAt: Date.now() - 60_000, - toWriteCount: 12, - phase: 'load-graph', - }, - }; - await saveMeta(storagePath, tampered); - const walGarbage = Buffer.alloc(8192, 0xab); - const shadowGarbage = Buffer.alloc(4096, 0xcd); - await writeFile(`${lbugPath}.wal`, walGarbage); - await writeFile(`${lbugPath}.shadow`, shadowGarbage); - - const logs: string[] = []; - // embeddingsNodeLimit: 1 (KTD9): the recovery runs force:true - // internally, and the seeded stats would otherwise route Phase 4 into - // a real embedder in CI — the 1-node cap suppresses generation while - // leaving the preserve/restore path fully live. On linux the - // wipe-and-restore vector-index seam then fires for real (statically - // linked VECTOR): a CREATE_VECTOR_INDEX over the restored rows is - // expected and harmless here. - const recovered = await runFullAnalysis( - repo.dbPath, - { skipAgentsMd: true, embeddingsNodeLimit: 1 }, - { onProgress: () => {}, onLog: (m) => logs.push(m) }, - ); - expect(recovered.alreadyUpToDate).toBeUndefined(); - - // Both sidecars were parked verbatim (renamed, never deleted) before - // any open could replay them… - expect(Buffer.compare(await readFile(`${lbugPath}.wal.dirty-recovery`), walGarbage)).toBe(0); - expect( - Buffer.compare(await readFile(`${lbugPath}.shadow.dirty-recovery`), shadowGarbage), - ).toBe(0); - const joinedLogs = logs.join('\n'); - expect(joinedLogs).toContain('Parked lbug.wal.dirty-recovery, lbug.shadow.dirty-recovery'); - - // …the run traversed the REAL pre-wipe preservation open — recovery's - // internal force on an embedded repo upgrades to regenerate mode, whose - // banner only prints when existingEmbeddingCount was read from the - // seeded stats and the cache-load path engaged… - expect(joinedLogs).toContain( - `--force on a repo with ${seededNodeIds.length} existing embeddings`, - ); - // …with generation itself cap-suppressed (no embedder in CI): - expect(joinedLogs).toContain('exceeds the 1-node safety cap'); - - // …and the rebuild completed into a clean index: dirty flag cleared, - // and the seeded embeddings survived the park → open → wipe → restore - // round-trip (the strongest signal the preservation open really ran: - // the DB was wiped, so these rows can only come from the cache load). - const after = await loadMeta(storagePath); - expect(after!.incrementalInProgress).toBeUndefined(); - expect(after!.stats?.embeddings).toBe(seededNodeIds.length); - } finally { - await repo.cleanup(); - } + it('restores seeded embeddings from a spill after wipe when the in-memory limit is 0', async () => { + await parkCrashSidecarsAndRebuild({ alwaysSpill: true }); }, 300_000); }); diff --git a/gitnexus/test/unit/run-analyze-fts-repair.test.ts b/gitnexus/test/unit/run-analyze-fts-repair.test.ts index 875e7dc47..2f0304ef4 100644 --- a/gitnexus/test/unit/run-analyze-fts-repair.test.ts +++ b/gitnexus/test/unit/run-analyze-fts-repair.test.ts @@ -7,7 +7,7 @@ import { saveMeta, type RepoMeta, } from '../../src/storage/repo-manager.js'; -import { EMBEDDING_DIMS } from '../../src/core/lbug/schema.js'; +import { EMBEDDING_DIMS, STALE_HASH_SENTINEL } from '../../src/core/lbug/schema.js'; import { getIndexIncompleteReasons } from '../../src/core/index-freshness.js'; import type { EmbeddingPipelineOptions, @@ -1225,6 +1225,313 @@ describe('runFullAnalysis wipe-and-restore vector-index stamp (tri-review 466951 await tmpRepo.cleanup(); } }); + + it('warns and continues without restore when loadCachedEmbeddings rejects (R11)', async () => { + const RESTORED_NODE_ID = 'Function:src/app.ts:handler:1'; + const stubNode = { + id: RESTORED_NODE_ID, + label: 'Function', + name: 'handler', + properties: { filePath: 'src/app.ts' }, + }; + const executeWithReusedStatement = vi.fn(async () => []); + vi.doMock('../../src/core/lbug/lbug-adapter.js', () => ({ + initLbug: vi.fn(async () => undefined), + loadGraphToLbug: vi.fn(async () => undefined), + getLbugStats: vi.fn(async () => ({ nodes: 2, edges: 0, communities: 0, processes: 0 })), + executeQuery: vi.fn(async () => []), + executeWithReusedStatement, + closeLbug: vi.fn(async () => undefined), + wipeLbugDbFiles: vi.fn(async () => undefined), + tryFlushWAL: vi.fn(async () => true), + loadCachedEmbeddings: vi.fn(async () => { + throw new Error('spill write failed'); + }), + deleteNodesForFile: vi.fn(async () => undefined), + deleteNodesForFiles: vi.fn(async () => undefined), + deleteAllCommunitiesAndProcesses: vi.fn(async () => undefined), + queryImporters: vi.fn(async () => []), + queryImportersBatch: vi.fn(async () => []), + loadFTSExtension: vi.fn(async () => false), + })); + vi.doMock('../../src/core/search/fts-indexes.js', () => ({ + initialiseSearchFTSStemmer: vi.fn(() => 'porter'), + createSearchFTSIndexes: vi.fn(async () => []), + verifySearchFTSIndexes: vi.fn(async () => []), + })); + vi.doMock('../../src/core/ingestion/pipeline.js', () => ({ + runPipelineFromRepo: vi.fn(async (repoPath: string) => ({ + repoPath, + totalFileCount: 1, + graph: { + forEachNode: (fn: (node: typeof stubNode) => void) => fn(stubNode), + getNode: (id: string) => (id === RESTORED_NODE_ID ? stubNode : undefined), + }, + })), + })); + vi.doMock('../../src/storage/repo-manager.js', async (importActual) => ({ + ...(await importActual()), + registerRepo: vi.fn(async () => 'cache-load-reject-repo'), + ensureGitNexusIgnored: vi.fn(async () => undefined), + })); + vi.doMock('../../src/core/embeddings/embedding-pipeline.js', async (importActual) => ({ + ...(await importActual()), + batchInsertEmbeddings: vi.fn(async () => { + throw new Error('restore must not run after a cache-load failure'); + }), + })); + + const tmpRepo = await createTempDir('gitnexus-run-analyze-cache-reject-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await fs.mkdir(storagePath, { recursive: true }); + await saveMeta(storagePath, { + repoPath: tmpRepo.dbPath, + lastCommit: '', + indexedAt: new Date().toISOString(), + stats: { embeddings: 1 }, + }); + + const logs: string[] = []; + const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); + await runFullAnalysis( + tmpRepo.dbPath, + { force: true, embeddingsNodeLimit: 1 }, + { onProgress: () => {}, onLog: (m) => logs.push(m) }, + ); + + expect(logs.some((m) => m.includes('Warning: could not load cached embeddings'))).toBe(true); + expect(logs.some((m) => m.includes('spill write failed'))).toBe(true); + expect( + executeWithReusedStatement.mock.calls.some((call) => + String(call[0]).includes('CREATE (e:CodeEmbedding'), + ), + ).toBe(false); + } finally { + await tmpRepo.cleanup(); + } + }); + + it('omits failed restore rows from the Phase 4 skip-set so they can be re-embedded', async () => { + const RESTORED_NODE_ID = 'Function:src/app.ts:handler:1'; + const CACHED_HASH = 'cached-stable-hash'; + const stubNode = { + id: RESTORED_NODE_ID, + label: 'Function', + name: 'handler', + properties: { filePath: 'src/app.ts' }, + }; + let existingEmbeddings: Map | undefined; + vi.doMock('../../src/core/lbug/lbug-adapter.js', () => ({ + initLbug: vi.fn(async () => undefined), + loadGraphToLbug: vi.fn(async () => undefined), + getLbugStats: vi.fn(async () => ({ nodes: 2, edges: 0, communities: 0, processes: 0 })), + executeQuery: vi.fn(async () => []), + executeWithReusedStatement: vi.fn(async () => []), + closeLbug: vi.fn(async () => undefined), + wipeLbugDbFiles: vi.fn(async () => undefined), + tryFlushWAL: vi.fn(async () => true), + loadCachedEmbeddings: vi.fn(async () => ({ + embeddingNodeIds: new Set([RESTORED_NODE_ID]), + embeddings: [ + { + nodeId: RESTORED_NODE_ID, + chunkIndex: 0, + startLine: 0, + endLine: 3, + embedding: new Array(EMBEDDING_DIMS).fill(0), + contentHash: CACHED_HASH, + }, + ], + })), + deleteNodesForFile: vi.fn(async () => undefined), + deleteNodesForFiles: vi.fn(async () => undefined), + deleteAllCommunitiesAndProcesses: vi.fn(async () => undefined), + queryImporters: vi.fn(async () => []), + queryImportersBatch: vi.fn(async () => []), + loadFTSExtension: vi.fn(async () => false), + })); + vi.doMock('../../src/core/search/fts-indexes.js', () => ({ + initialiseSearchFTSStemmer: vi.fn(() => 'porter'), + createSearchFTSIndexes: vi.fn(async () => []), + verifySearchFTSIndexes: vi.fn(async () => []), + })); + vi.doMock('../../src/core/ingestion/pipeline.js', () => ({ + runPipelineFromRepo: vi.fn(async (repoPath: string) => ({ + repoPath, + totalFileCount: 1, + graph: { + forEachNode: (fn: (node: typeof stubNode) => void) => fn(stubNode), + getNode: (id: string) => (id === RESTORED_NODE_ID ? stubNode : undefined), + }, + })), + })); + vi.doMock('../../src/storage/repo-manager.js', async (importActual) => ({ + ...(await importActual()), + registerRepo: vi.fn(async () => 'ktd7-skip-set-repo'), + ensureGitNexusIgnored: vi.fn(async () => undefined), + })); + const runEmbeddingPipeline = vi.fn( + async ( + _executeQuery: unknown, + _executeWithReusedStatement: unknown, + _onProgress: unknown, + _config: unknown, + _cachedNodeIds: unknown, + embeddings: Map | undefined, + ) => { + existingEmbeddings = embeddings; + return { + nodesProcessed: 0, + chunksProcessed: 0, + vectorIndexReady: false, + semanticMode: 'exact-scan' as const, + failedNodeIds: [], + }; + }, + ); + vi.doMock('../../src/core/embeddings/embedding-pipeline.js', () => ({ + runEmbeddingPipeline, + buildVectorIndex: vi.fn(async () => false), + batchInsertEmbeddings: vi.fn(async () => { + throw new Error('nth batch insert failed'); + }), + })); + + const tmpRepo = await createTempDir('gitnexus-run-analyze-ktd7-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await fs.mkdir(storagePath, { recursive: true }); + await saveMeta(storagePath, { + repoPath: tmpRepo.dbPath, + lastCommit: '', + indexedAt: new Date().toISOString(), + stats: { embeddings: 1 }, + }); + + const logs: string[] = []; + const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); + await runFullAnalysis( + tmpRepo.dbPath, + { force: true }, + { onProgress: () => {}, onLog: (m) => logs.push(m) }, + ); + + expect(runEmbeddingPipeline).toHaveBeenCalled(); + expect(existingEmbeddings?.get(RESTORED_NODE_ID)).toBe(STALE_HASH_SENTINEL); + expect(logs.some((m) => m.includes('Warning: could not restore'))).toBe(true); + } finally { + await tmpRepo.cleanup(); + } + }); + + it('marks a node stale when only some of its restore batches succeed', async () => { + const RESTORED_NODE_ID = 'Function:src/app.ts:handler:1'; + const CACHED_HASH = 'cached-stable-hash'; + const stubNode = { + id: RESTORED_NODE_ID, + label: 'Function', + name: 'handler', + properties: { filePath: 'src/app.ts' }, + }; + const embeddings = Array.from({ length: 201 }, (_, chunkIndex) => ({ + nodeId: RESTORED_NODE_ID, + chunkIndex, + startLine: chunkIndex, + endLine: chunkIndex + 1, + embedding: new Array(EMBEDDING_DIMS).fill(0), + contentHash: CACHED_HASH, + })); + let existingEmbeddings: Map | undefined; + vi.doMock('../../src/core/lbug/lbug-adapter.js', () => ({ + initLbug: vi.fn(async () => undefined), + loadGraphToLbug: vi.fn(async () => undefined), + getLbugStats: vi.fn(async () => ({ nodes: 2, edges: 0, communities: 0, processes: 0 })), + executeQuery: vi.fn(async () => []), + executeWithReusedStatement: vi.fn(async () => []), + closeLbug: vi.fn(async () => undefined), + wipeLbugDbFiles: vi.fn(async () => undefined), + tryFlushWAL: vi.fn(async () => true), + loadCachedEmbeddings: vi.fn(async () => ({ + embeddingNodeIds: new Set([RESTORED_NODE_ID]), + embeddings, + })), + deleteNodesForFile: vi.fn(async () => undefined), + deleteNodesForFiles: vi.fn(async () => undefined), + deleteAllCommunitiesAndProcesses: vi.fn(async () => undefined), + queryImporters: vi.fn(async () => []), + queryImportersBatch: vi.fn(async () => []), + loadFTSExtension: vi.fn(async () => false), + })); + vi.doMock('../../src/core/search/fts-indexes.js', () => ({ + initialiseSearchFTSStemmer: vi.fn(() => 'porter'), + createSearchFTSIndexes: vi.fn(async () => []), + verifySearchFTSIndexes: vi.fn(async () => []), + })); + vi.doMock('../../src/core/ingestion/pipeline.js', () => ({ + runPipelineFromRepo: vi.fn(async (repoPath: string) => ({ + repoPath, + totalFileCount: 1, + graph: { + forEachNode: (fn: (node: typeof stubNode) => void) => fn(stubNode), + getNode: (id: string) => (id === RESTORED_NODE_ID ? stubNode : undefined), + }, + })), + })); + vi.doMock('../../src/storage/repo-manager.js', async (importActual) => ({ + ...(await importActual()), + registerRepo: vi.fn(async () => 'partial-restore-skip-set-repo'), + ensureGitNexusIgnored: vi.fn(async () => undefined), + })); + const runEmbeddingPipeline = vi.fn( + async ( + _executeQuery: unknown, + _executeWithReusedStatement: unknown, + _onProgress: unknown, + _config: unknown, + _cachedNodeIds: unknown, + embeddings: Map | undefined, + ) => { + existingEmbeddings = embeddings; + return { + nodesProcessed: 0, + chunksProcessed: 0, + vectorIndexReady: false, + semanticMode: 'exact-scan' as const, + failedNodeIds: [], + }; + }, + ); + let inserts = 0; + vi.doMock('../../src/core/embeddings/embedding-pipeline.js', () => ({ + runEmbeddingPipeline, + buildVectorIndex: vi.fn(async () => false), + batchInsertEmbeddings: vi.fn(async () => { + inserts += 1; + if (inserts > 1) throw new Error('second batch insert failed'); + }), + })); + + const tmpRepo = await createTempDir('gitnexus-run-analyze-partial-restore-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await fs.mkdir(storagePath, { recursive: true }); + await saveMeta(storagePath, { + repoPath: tmpRepo.dbPath, + lastCommit: '', + indexedAt: new Date().toISOString(), + stats: { embeddings: 201 }, + }); + + const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); + await runFullAnalysis(tmpRepo.dbPath, { force: true }, { onProgress: () => {} }); + + expect(inserts).toBe(2); + expect(existingEmbeddings?.get(RESTORED_NODE_ID)).toBe(STALE_HASH_SENTINEL); + } finally { + await tmpRepo.cleanup(); + } + }); }); /** diff --git a/gitnexus/vitest.config.ts b/gitnexus/vitest.config.ts index 4b2625805..1239862a3 100644 --- a/gitnexus/vitest.config.ts +++ b/gitnexus/vitest.config.ts @@ -100,6 +100,7 @@ export default defineConfig({ 'test/integration/analyze-wal-checkpoint-failure.test.ts', 'test/integration/lbug-non-ascii-path.test.ts', 'test/integration/lbug-conn-serialization.test.ts', + 'test/integration/load-cached-embeddings-spill.test.ts', 'test/integration/group/manifest-resolve-symbol-2325.test.ts', 'test/integration/group/manifest-synthetic-impact-lbug.test.ts', 'test/integration/group/http-route-resolve-symbol.test.ts', @@ -179,6 +180,7 @@ export default defineConfig({ 'test/integration/analyze-wal-checkpoint-failure.test.ts', 'test/integration/lbug-non-ascii-path.test.ts', 'test/integration/lbug-conn-serialization.test.ts', + 'test/integration/load-cached-embeddings-spill.test.ts', 'test/integration/group/manifest-resolve-symbol-2325.test.ts', 'test/integration/group/manifest-synthetic-impact-lbug.test.ts', 'test/integration/group/http-route-resolve-symbol.test.ts',