fix(embeddings): stop Caching embeddings OOM on large incremental analyze (#3310)

* fix(embeddings): spill cached vectors to a Float32 temp file

Keep restore metadata in RAM and write embeddings once the in-memory
row limit is exceeded so incremental analyze can survive large caches
without a full-table number[] heap (#3306).

Co-authored-by: Cursor <cursoragent@cursor.com>

* fix(lbug): stream CodeEmbedding cache under the connection lock

Spill vectors once the in-memory limit is crossed and fail the load
instead of adopting an empty snapshot, so incremental analyze cannot
OOM or quietly drop the restore cache (#3306).

Co-authored-by: Cursor <cursoragent@cursor.com>

* fix(analyze): restore cached embeddings from a streamed spill snapshot

Hold row metadata across wipe, materialize 200-row batches, and treat
cache-load failures as warn-and-continue so incremental analyze can
preserve vectors without a full-table heap (#3306).

Co-authored-by: Cursor <cursoragent@cursor.com>

* Address PR review feedback (#3310)

Loop spill writes until the full vector lands, keep materialize failures out of the insert catch and the Phase 4 hash skip-set, and assert spilled restore subsets by node id instead of scan order.

Co-authored-by: Cursor <cursoragent@cursor.com>

* Address PR review feedback (#3310)

Discard only this analyze run's embedding spills so a concurrent analyze on another index keeps its restore file, and isolate the default in-memory limit test from inherited env.

Co-authored-by: Cursor <cursoragent@cursor.com>

* Address PR review feedback (#3310)

Mark a node stale when any restore batch fails so leftover chunks are deleted and rembedded, and exercise a full-length bad-magic spill header.

---------

Co-authored-by: Gergo Magyar <gergomagyar0@gmail.com>
Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
Gergő Magyar 2026-09-17 16:18:29 +01:00 • committed by GitHub
parent d2a43e33df
commit a2e1710ac0
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
8 changed files with 1493 additions and 218 deletions

View file

@ -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<string>;
/** 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<string>;
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<string>;
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<ArrayBufferView, DataView> & { 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<unknown>)[Symbol.iterator] === 'function'
) {
const arr = Array.isArray(embedding)
? (embedding as unknown[])
: Array.from(embedding as Iterable<unknown>);
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<number>(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<string>();
const spillScope = new AsyncLocalStorage<Set<string>>();
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<T>(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<string, unknown> | unknown[],
hasContentHash: boolean,
): void {
const rec = row as Record<string, unknown> & 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;
});
}

View file

@ -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<string>;
embeddings: CachedEmbedding[];
}> => {
export const loadCachedEmbeddings = async (
options?: LoadCachedEmbeddingsOptions,
): Promise<CachedEmbeddingsSnapshot> => {
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<string>();
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 };
});
};

View file

@ -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<string>();
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<string>();
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<string>();
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<nodeId, contentHash> from cached embeddings for incremental mode
let existingEmbeddings: Map<string, string> | undefined;
if (cachedEmbeddingNodeIds.size > 0) {
if (cachedSnapshot.embeddingNodeIds.size > 0) {
existingEmbeddings = new Map<string, string>();
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`.

View file

@ -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,
);
});
});

View file

@ -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/);
});
});

View file

@ -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);
});

View file

@ -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<typeof import('../../src/storage/repo-manager.js')>()),
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<typeof import('../../src/core/embeddings/embedding-pipeline.js')>()),
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<string, string> | 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<typeof import('../../src/storage/repo-manager.js')>()),
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<string, string> | 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<string, string> | 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<typeof import('../../src/storage/repo-manager.js')>()),
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<string, string> | 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();
}
});
});
/**

View file

@ -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',