diff --git a/docker-server.test.mjs b/docker-server.test.mjs index 6d2a9f6c2..309edf3b5 100644 --- a/docker-server.test.mjs +++ b/docker-server.test.mjs @@ -331,8 +331,8 @@ const respondOk = (_req, res) => { // // upstream request handler, replaceable mid-test via `ctx.handler`; // null points the proxy at a port nothing ever listens on -// listenAfterMs bind the upstream this late, so the first attempt(s) hit -// ECONNREFUSED (a single-instance restart window) +// listenAfterMs bind the upstream this late after the first refused attempt +// (a single-instance restart window) // schemeless drop http:// from GITNEXUS_UPSTREAM_URL, the way Render's // `fromService: { property: hostport }` yields it // env extra environment for docker-server.mjs @@ -366,18 +366,15 @@ async function withProxy( }) : null; - // A late (or never) bind needs its port reserved up front; otherwise let the - // OS assign one at listen time. - const upstreamPort = - server && listenAfterMs === 0 - ? await new Promise((r) => server.listen(0, '127.0.0.1', () => r(server.address().port))) - : await getFreePort(); - const bindTimer = - server && listenAfterMs > 0 - ? setTimeout(() => server.listen(upstreamPort, '127.0.0.1'), listenAfterMs) - : null; - + // Keep the upstream port bound until the proxy port is chosen. Releasing it + // sooner lets the OS assign both services the same port and proxy to itself. + const reservation = server ?? createServer(); + const upstreamPort = await new Promise((r) => + reservation.listen(0, '127.0.0.1', () => r(reservation.address().port)), + ); const port = await getFreePort(); + let bindTimer = null; + const target = `127.0.0.1:${upstreamPort}`; const proc = spawnServerWithEnv(dir, port, { GITNEXUS_UPSTREAM_URL: schemeless ? target : `http://${target}`, @@ -390,18 +387,25 @@ async function withProxy( ...env, }); proc.stderr.setEncoding('utf8'); - proc.stderr.on('data', (chunk) => { + const collectStderr = (chunk) => { ctx.stderr += chunk; - }); + // Process startup must not consume the restart window or skip the retry. + if (server && listenAfterMs > 0 && !bindTimer && ctx.stderr.includes('ECONNREFUSED; retry')) { + bindTimer = setTimeout(() => server.listen(upstreamPort, '127.0.0.1'), listenAfterMs); + } + }; + proc.stderr.on('data', collectStderr); try { await waitForServer(port); + if (!server || listenAfterMs > 0) await new Promise((r) => reservation.close(r)); await fn(port, ctx); } finally { + proc.stderr.off('data', collectStderr); if (bindTimer) clearTimeout(bindTimer); await killAndWait(proc); - if (server?.listening) { - server.closeAllConnections?.(); - await new Promise((r) => server.close(r)); + if (reservation.listening) { + reservation.closeAllConnections?.(); + await new Promise((r) => reservation.close(r)); } await rm(dir, { recursive: true, force: true }); } @@ -661,8 +665,8 @@ it('returns 502 when the upstream is unreachable', async () => { // -- Connection-retry across an upstream restart window --------------------- // -// `listenAfterMs: 400` binds the upstream late, so the first attempt hits -// ECONNREFUSED and must be retried — a single-instance restart. The default 3 +// `listenAfterMs: 400` binds the upstream 400ms after the first ECONNREFUSED, +// so the request must be retried — a single-instance restart. The default 3 // attempts (backoff 250ms, 500ms) span ~750ms, so a retry lands after the bind. it('retries a connection-refused POST and succeeds once the upstream is up', async () => { @@ -676,6 +680,7 @@ it('retries a connection-refused POST and succeeds once the upstream is up', asy assert.match(res.body, /"ok":true/); assert.equal(ctx.calls, 1, 'upstream must run the job exactly once (no double-execute)'); assert.equal(ctx.body, '{"repo":"x"}', 'buffered body replayed intact'); + assert.match(ctx.stderr, /ECONNREFUSED; retry/, 'the restart gap must exercise a retry'); }); }); diff --git a/gitnexus/src/cli/embeddings-sync.ts b/gitnexus/src/cli/embeddings-sync.ts index 9848ef7c4..59cf9270a 100644 --- a/gitnexus/src/cli/embeddings-sync.ts +++ b/gitnexus/src/cli/embeddings-sync.ts @@ -4,7 +4,11 @@ import { LBUG_DIRECTORY } from '../storage/storage-constants.js'; import path from 'node:path'; import { cliInfo } from './cli-message.js'; import { getGitRoot } from '../storage/git.js'; -import { acquireIndexLock, requireExclusiveIndexLock } from '../storage/index-lock.js'; +import { + acquireIndexLock, + requireExclusiveIndexLock, + sweepStagingArtifacts, +} from '../storage/index-lock.js'; import { getStoragePaths, loadMeta, saveMeta } from '../storage/repo-manager.js'; import { closeLbug, @@ -50,12 +54,22 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise = // Writes go to the slot's own graph. A shared-store checkout that reads an // immutable commit graph (#3352) takes a private copy first. const lbugPath = path.join(metaDir, LBUG_DIRECTORY); - const lock = await acquireIndexLock(metaDir); + const lock = await acquireIndexLock(metaDir, { sweep: false }); try { requireExclusiveIndexLock( lock, `Cannot acquire the index lock at ${metaDir}; refusing an unlocked embeddings sync.`, ); + // Sync writes the published graph and cannot recover a staged generation. + // Reject even malformed receipts before sweeping staging files or writing. + const recoveryCheckpoint = (await loadMeta(metaDir))?.embeddingCheckpoint; + if (recoveryCheckpoint && Object.hasOwn(recoveryCheckpoint, 'recovery')) { + throw new Error( + 'Cannot sync embeddings: the index checkpoint references staged embeddings. ' + + 'Run `gitnexus analyze` to recover them first.', + ); + } + sweepStagingArtifacts(metaDir); if (!(await ensurePrivateSharedGraph(metaDir, (m) => console.log(` ${m}`)))) { throw new Error('The shared graph this checkout reads is gone. Run gitnexus analyze first.'); } diff --git a/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts b/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts new file mode 100644 index 000000000..b1fb2cb73 --- /dev/null +++ b/gitnexus/src/core/embeddings/staged-embedding-recovery-child.ts @@ -0,0 +1,181 @@ +/** Isolated strict native reader. Never import this entrypoint into analyze. */ +import fs from 'node:fs'; +import path from 'node:path'; +import lbug from '@ladybugdb/core'; +import { createLbugDatabase, toNativeSafePath } from '../lbug/lbug-config.js'; +import { FAMILY_SUFFIXES } from '../../storage/embedding-recovery.js'; +import { + abortCachedEmbeddingsBuilder, + createCachedEmbeddingsBuilder, + finalizeCachedEmbeddingsSnapshot, + ingestCachedEmbeddingRow, +} from './embedding-restore-spill.js'; +import type { StagedEmbeddingExport } from './staged-embedding-recovery.js'; + +/** Native replay and close may checkpoint; only give them disposable copies. */ +function copyRecoveryFamily(dbPath: string, exportDir: string): string { + const replayDir = fs.mkdtempSync(path.join(exportDir, 'replay-')); + const replayPath = path.join(replayDir, path.basename(dbPath)); + const noFollow = fs.constants.O_NOFOLLOW ?? 0; + const flags = fs.constants.O_RDONLY | noFollow | (fs.constants.O_NONBLOCK ?? 0); + const buffer = Buffer.allocUnsafe(1024 * 1024); + for (const suffix of FAMILY_SUFFIXES) { + const sourcePath = dbPath + suffix; + let entry: fs.BigIntStats; + try { + entry = fs.lstatSync(sourcePath, { bigint: true }); + } catch (error) { + if (suffix && (error as NodeJS.ErrnoException).code === 'ENOENT') continue; + throw error; + } + if (!entry.isFile()) throw new Error('staged embedding family is not a regular file'); + const source = fs.openSync(sourcePath, flags); + let destination: number | undefined; + try { + const opened = fs.fstatSync(source, { bigint: true }); + const current = fs.lstatSync(sourcePath, { bigint: true }); + if ( + !opened.isFile() || + !current.isFile() || + (noFollow === 0 && opened.ino === 0n) || + opened.dev !== entry.dev || + opened.ino !== entry.ino || + opened.dev !== current.dev || + opened.ino !== current.ino + ) { + throw new Error('staged embedding family changed while opening'); + } + destination = fs.openSync(replayPath + suffix, 'wx', 0o600); + let copied = 0n; + for (;;) { + const bytesRead = fs.readSync(source, buffer, 0, buffer.length, null); + if (bytesRead === 0) break; + fs.writeFileSync(destination, buffer.subarray(0, bytesRead)); + copied += BigInt(bytesRead); + } + const after = fs.fstatSync(source, { bigint: true }); + if ( + copied !== opened.size || + after.size !== opened.size || + after.mtimeNs !== opened.mtimeNs || + after.ctimeNs !== opened.ctimeNs + ) { + throw new Error('staged embedding family changed while copying'); + } + } finally { + try { + if (destination !== undefined) fs.closeSync(destination); + } finally { + fs.closeSync(source); + } + } + } + return replayPath; +} + +async function extract(): Promise { + const [dbPath, exportDir, dimensionsArg] = process.argv.slice(2); + const dimensions = Number(dimensionsArg); + if (!dbPath || !exportDir || !Number.isInteger(dimensions) || dimensions <= 0) { + throw new Error('invalid staged embedding extraction arguments'); + } + // The parent owns exportDir and reclaims it even after killing this child. + const replayPath = copyRecoveryFamily(dbPath, exportDir); + const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0, spillDir: exportDir }); + const rejectedNodeIds = new Set(); + let db: lbug.Database | undefined; + let conn: lbug.Connection | undefined; + try { + // Avoid openLbugConnection's test-fixture lock sweep: a recovery source must + // never have its WAL removed, even when an external slot resembles a fixture. + db = createLbugDatabase(lbug, toNativeSafePath(replayPath), { throwOnWalReplayFailure: true }); + conn = new lbug.Connection(db); + const queried = await conn.query( + 'MATCH (e:CodeEmbedding) 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', + ); + const results = Array.isArray(queried) ? queried : [queried]; + try { + if (results.length !== 1) throw new Error('unexpected staged embedding query result'); + const result = results[0]; + while (await result.hasNext()) { + const raw = await result.getNext(); + const rec = raw as Record & unknown[]; + const nodeId = rec.nodeId ?? rec[0]; + if (typeof nodeId !== 'string' || !nodeId) + throw new Error('invalid staged embedding node id'); + const chunkIndex = rec.chunkIndex ?? rec[1]; + const startLine = rec.startLine ?? rec[2]; + const endLine = rec.endLine ?? rec[3]; + const embedding = rec.embedding ?? rec[4]; + const contentHash = rec.contentHash ?? rec[5]; + const vector = + Array.isArray(embedding) || + (ArrayBuffer.isView(embedding) && !(embedding instanceof DataView)) + ? Array.from(embedding as ArrayLike) + : undefined; + if ( + !Number.isInteger(chunkIndex) || + Number(chunkIndex) < 0 || + !Number.isInteger(startLine) || + Number(startLine) < 0 || + !Number.isInteger(endLine) || + Number(endLine) < Number(startLine) || + typeof contentHash !== 'string' || + !contentHash || + !vector || + vector.length !== dimensions || + vector.some( + (value) => + typeof value !== 'number' || + !Number.isFinite(value) || + !Number.isFinite(Math.fround(value)), + ) + ) { + rejectedNodeIds.add(nodeId); + continue; + } + ingestCachedEmbeddingRow( + builder, + { nodeId, chunkIndex, startLine, endLine, embedding: vector, contentHash }, + true, + ); + } + } finally { + for (const result of results) await result.close(); + } + // Both closes must succeed. Suppressed native teardown errors are unsafe. + await conn.close(); + await db.close(); + const snapshot = finalizeCachedEmbeddingsSnapshot(builder); + if (snapshot.spill) fs.renameSync(snapshot.spill.path, path.join(exportDir, 'vectors.bin')); + const manifest: StagedEmbeddingExport = { + version: 1, + dimensions, + rows: snapshot.rows, + rejectedNodeIds: [...rejectedNodeIds], + }; + fs.writeFileSync(path.join(exportDir, 'manifest.json'), JSON.stringify(manifest), { + flag: 'wx', + mode: 0o600, + }); + } catch (err) { + abortCachedEmbeddingsBuilder(builder); + // Cleanup is best effort on a rejected source, never used to approve output. + try { + await conn?.close(); + } catch { + /* rejected */ + } + try { + await db?.close(); + } catch { + /* rejected */ + } + throw err; + } +} + +extract().catch((err: unknown) => { + process.stderr.write(`${err instanceof Error ? err.message : String(err)}\n`); + process.exitCode = 1; +}); diff --git a/gitnexus/src/core/embeddings/staged-embedding-recovery.ts b/gitnexus/src/core/embeddings/staged-embedding-recovery.ts new file mode 100644 index 000000000..b295f4beb --- /dev/null +++ b/gitnexus/src/core/embeddings/staged-embedding-recovery.ts @@ -0,0 +1,244 @@ +/** Recover paid embedding rows without opening an interrupted native DB in analyze. */ +import { spawn } from 'node:child_process'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { + abortCachedEmbeddingsBuilder, + createCachedEmbeddingsBuilder, + EmbeddingSpillReader, + emptyCachedEmbeddingsSnapshot, + finalizeCachedEmbeddingsSnapshot, + ingestCachedEmbeddingRow, + materializeCachedEmbeddings, + type CachedEmbeddingMeta, + type CachedEmbeddingsSnapshot, +} from './embedding-restore-spill.js'; + +export interface StagedEmbeddingRecoveryOptions { + dimensions: number; + /** Active checkpoint window and any incomplete inherited restore groups. */ + excludedNodeIds?: Iterable; + timeoutMs?: number; +} + +export interface StagedEmbeddingExport { + version: 1; + dimensions: number; + rows: CachedEmbeddingMeta[]; + rejectedNodeIds: string[]; +} + +const RESTORE_BATCH_SIZE = 200; +const DEFAULT_EXTRACTION_TIMEOUT_MS = 120_000; + +/** A contiguous prefix is safe only when its node is outside every unsafe window. */ +export function validateRecoveredNodeGroups( + rows: readonly CachedEmbeddingMeta[], + excludedNodeIds: ReadonlySet = new Set(), +): Set { + const groups = new Map; invalid: boolean }>(); + for (const row of rows) { + let group = groups.get(row.nodeId); + if (!group) { + group = { hash: row.contentHash, ordinals: new Set(), invalid: false }; + groups.set(row.nodeId, group); + } + if ( + typeof row.nodeId !== 'string' || + !row.nodeId || + typeof row.contentHash !== 'string' || + !row.contentHash || + row.contentHash !== group.hash || + !Number.isInteger(row.chunkIndex) || + row.chunkIndex < 0 || + group.ordinals.has(row.chunkIndex) || + !Number.isInteger(row.startLine) || + row.startLine < 0 || + !Number.isInteger(row.endLine) || + row.endLine < row.startLine + ) + group.invalid = true; + group.ordinals.add(row.chunkIndex); + } + const accepted = new Set(); + for (const [nodeId, group] of groups) { + if (group.invalid || excludedNodeIds.has(nodeId)) continue; + // Unique ordinals with max n-1 and zero present have no holes. + if (!group.ordinals.has(0)) continue; + if ([...group.ordinals].some((ordinal) => ordinal >= group.ordinals.size)) continue; + accepted.add(nodeId); + } + return accepted; +} + +/** Whole recovered nodes replace whole published groups; vectors stay in bounded batches. */ +export function mergeRecoveredEmbeddings( + live: CachedEmbeddingsSnapshot, + recovered: CachedEmbeddingsSnapshot, +): CachedEmbeddingsSnapshot { + const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0 }); + try { + appendSnapshotRows( + live, + live.rows.filter((row) => !recovered.embeddingNodeIds.has(row.nodeId)), + builder, + ); + appendSnapshotRows(recovered, recovered.rows, builder); + return finalizeCachedEmbeddingsSnapshot(builder); + } catch (err) { + abortCachedEmbeddingsBuilder(builder); + throw err; + } +} + +function appendSnapshotRows( + snapshot: CachedEmbeddingsSnapshot, + rows: readonly CachedEmbeddingMeta[], + builder: ReturnType, +): void { + const reader = snapshot.spill ? new EmbeddingSpillReader(snapshot.spill) : undefined; + try { + for (let i = 0; i < rows.length; i += RESTORE_BATCH_SIZE) { + const batch = materializeCachedEmbeddings( + snapshot, + rows.slice(i, i + RESTORE_BATCH_SIZE), + reader, + ); + for (const row of batch) + ingestCachedEmbeddingRow(builder, row as unknown as Record, true); + } + } finally { + reader?.close(); + } +} + +/** + * Caller must validate checkpoint identity, schema, exact generation, and hold the + * index lock. A subprocess contains native WAL replay/query/destructor failures. + * A failed strict open is never retried with WAL removed or validation disabled. + */ +export async function recoverStagedEmbeddings( + dbPath: string, + options: StagedEmbeddingRecoveryOptions, +): Promise { + if (!Number.isInteger(options.dimensions) || options.dimensions <= 0) { + throw new Error('invalid staged embedding dimensions'); + } + const dbStat = fs.lstatSync(dbPath); + if (!dbStat.isFile() || dbStat.isSymbolicLink()) + throw new Error('staged embedding DB is not a regular file'); + const exportDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-export-')); + try { + await runExtractionChild(dbPath, exportDir, options); + const manifestPath = path.join(exportDir, 'manifest.json'); + const manifest = JSON.parse(fs.readFileSync(manifestPath, 'utf8')) as StagedEmbeddingExport; + if ( + manifest.version !== 1 || + manifest.dimensions !== options.dimensions || + !Array.isArray(manifest.rows) || + !Array.isArray(manifest.rejectedNodeIds) || + manifest.rejectedNodeIds.some((id) => typeof id !== 'string') || + manifest.rows.some((row, i) => !row || row.vectorIndex !== i) + ) + throw new Error('invalid staged embedding export manifest'); + if (manifest.rows.length === 0) return emptyCachedEmbeddingsSnapshot(); + const spill = { + path: path.join(exportDir, 'vectors.bin'), + dims: options.dimensions, + rowCount: manifest.rows.length, + }; + const vectorStat = fs.lstatSync(spill.path); + if ( + !vectorStat.isFile() || + vectorStat.isSymbolicLink() || + vectorStat.size !== 12 + spill.rowCount * spill.dims * 4 + ) { + throw new Error('invalid staged embedding export size'); + } + const excluded = new Set(options.excludedNodeIds ?? []); + for (const nodeId of manifest.rejectedNodeIds) excluded.add(nodeId); + const accepted = validateRecoveredNodeGroups(manifest.rows, excluded); + const exported: CachedEmbeddingsSnapshot = { + rows: manifest.rows, + embeddings: [], + embeddingNodeIds: accepted, + spill, + }; + const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0 }); + const reader = new EmbeddingSpillReader(spill); + try { + const rows = manifest.rows.filter((row) => accepted.has(row.nodeId)); + for (let i = 0; i < rows.length; i += RESTORE_BATCH_SIZE) { + for (const row of materializeCachedEmbeddings( + exported, + rows.slice(i, i + RESTORE_BATCH_SIZE), + reader, + )) { + if ( + row.embedding.length !== options.dimensions || + row.embedding.some((value) => !Number.isFinite(value)) + ) { + throw new Error('invalid staged embedding vector'); + } + ingestCachedEmbeddingRow(builder, row as unknown as Record, true); + } + } + return finalizeCachedEmbeddingsSnapshot(builder); + } catch (err) { + abortCachedEmbeddingsBuilder(builder); + throw err; + } finally { + reader.close(); + } + } finally { + fs.rmSync(exportDir, { recursive: true, force: true }); + } +} + +function runExtractionChild( + dbPath: string, + exportDir: string, + options: StagedEmbeddingRecoveryOptions, +): Promise { + const compiledPath = fileURLToPath( + new URL('./staged-embedding-recovery-child.js', import.meta.url), + ); + const sourcePath = compiledPath.replace(/\.js$/, '.ts'); + const childPath = fs.existsSync(compiledPath) ? compiledPath : sourcePath; + const args = childPath.endsWith('.ts') + ? ['--import', import.meta.resolve('tsx'), childPath] + : [childPath]; + args.push(dbPath, exportDir, String(options.dimensions)); + return new Promise((resolve, reject) => { + const child = spawn(process.execPath, args, { + stdio: ['ignore', 'ignore', 'pipe'], + windowsHide: true, + }); + let stderr = ''; + let timedOut = false; + child.stderr.on('data', (chunk: Buffer) => { + if (stderr.length < 4096) stderr += chunk.toString().slice(0, 4096 - stderr.length); + }); + const timer = setTimeout(() => { + timedOut = true; + // A process boundary is safe to kill even while native code is executing. + child.kill('SIGKILL'); + }, options.timeoutMs ?? DEFAULT_EXTRACTION_TIMEOUT_MS); + child.once('error', (err) => { + clearTimeout(timer); + reject(err); + }); + child.once('close', (code, signal) => { + clearTimeout(timer); + if (code === 0 && !signal && !timedOut) resolve(); + else + reject( + new Error( + `staged embedding extraction failed${timedOut ? ' (timeout)' : ` (${signal ?? code})`}${stderr ? `: ${stderr.trim()}` : ''}`, + ), + ); + }); + }); +} diff --git a/gitnexus/src/core/run-analyze.ts b/gitnexus/src/core/run-analyze.ts index af51bb6ea..5f43df3e2 100644 --- a/gitnexus/src/core/run-analyze.ts +++ b/gitnexus/src/core/run-analyze.ts @@ -36,7 +36,12 @@ import fs from 'fs/promises'; import { constants as fsConstants, existsSync } from 'node:fs'; import { randomUUID } from 'node:crypto'; import { retryRename } from '../storage/fs-atomic.js'; -import { acquireIndexLock, requireExclusiveIndexLock } from '../storage/index-lock.js'; +import { + acquireIndexLock, + requireExclusiveIndexLock, + sweepStagingArtifacts, +} from '../storage/index-lock.js'; +import { resolveEmbeddingRecovery } from '../storage/embedding-recovery.js'; import { invalidateNodeWorkspacePackages } from './ingestion/import-resolvers/node-workspace-packages.js'; import { logNameFallbackSummary, @@ -1287,6 +1292,8 @@ export async function runFullAnalysis( const log = (msg: string) => callbacks.onLog?.(stripControlCharacters(msg)); const acquireOpts = { log, + // Resolve and validate the canonical slot under the lock before cleanup. + sweep: false, onWaitStart: () => callbacks.onProgress('lock', 0, 'Waiting for another analyze to finish on this index…'), }; @@ -1345,6 +1352,7 @@ export async function runFullAnalysis( } const flatShared = writeTarget.placement.branch ? undefined : writeTarget.sharedStore; if (flatShared) await seedSharedSlot(flatShared, repoPath, log); + sweepStagingArtifacts(writeTarget.metaDir, log); const slotToLeave = options.noShare ? await optedInSlotToLeave(repoPath) : undefined; const result = await runFullAnalysisInner( repoPath, @@ -1842,7 +1850,13 @@ async function runFullAnalysisInner( decision = decideEmbeddingResume(checkpoint, embeddingIdentityForRun, resumeOptions); } if (decision.action === 'abort') throw new Error(decision.error); - log(decision.log); + log( + decision.action === 'resume' && checkpoint.recovery + ? `Previous analyze recorded an embedding checkpoint (${checkpoint.nodesProcessed}/` + + `${checkpoint.totalNodes} nodes); validating staged vectors before retrying ` + + `${decision.pendingNodeIds.size} pending node(s).` + : decision.log, + ); if (options.dropEmbeddings) { // --drop-embeddings has always implied a rebuild here; the decision only // covers the marker. @@ -2671,6 +2685,74 @@ async function runFullAnalysisInner( } } + // A checkpoint's pending decision and its paid, complete vectors are + // independent: --force discards the former but can still reuse the latter. + // Select only the explicitly referenced generation, never an orphan by age. + const stagedRecovery = resolveEmbeddingRecovery(metaDir, existingMeta?.embeddingCheckpoint); + const stagedCheckpoint = existingMeta?.embeddingCheckpoint; + const inheritedUnsafeNodeIds = new Set([ + ...pendingEmbeddingNodeIds, + ...(stagedRecovery?.unsafeNodeIds ?? []), + ]); + if (shouldLoadCache && !options.dropEmbeddings && stagedRecovery && stagedCheckpoint) { + if (!embeddingIdentityForRun) { + const { resolveEmbeddingIdentity } = await import('./embeddings/embedding-identity.js'); + embeddingIdentityForRun = resolveEmbeddingIdentity(); + } + const marker = stagedCheckpoint; + const matchesIdentity = + marker.model === embeddingIdentityForRun.model && + marker.dimensions === embeddingIdentityForRun.dimensions && + marker.provider === embeddingIdentityForRun.provider; + if (matchesIdentity && stagedRecovery.schemaFingerprint === SCHEMA_FINGERPRINT) { + // A force-discarded pending decision must not turn a known incomplete + // inherited group into a reusable cache merely because its hash matches. + if (inheritedUnsafeNodeIds.size > 0) { + const rows = cachedSnapshot.rows.filter((row) => !inheritedUnsafeNodeIds.has(row.nodeId)); + cachedSnapshot = { + ...cachedSnapshot, + rows, + embeddingNodeIds: new Set(rows.map((row) => row.nodeId)), + }; + } + let recovered: CachedEmbeddingsSnapshot | undefined; + try { + const { recoverStagedEmbeddings, mergeRecoveredEmbeddings } = + await import('./embeddings/staged-embedding-recovery.js'); + recovered = await recoverStagedEmbeddings(stagedRecovery.dbPath, { + dimensions: embeddingIdentityForRun.dimensions, + excludedNodeIds: stagedRecovery.unsafeNodeIds, + }); + const liveCacheDims = snapshotEmbeddingDims(cachedSnapshot); + if (liveCacheDims !== undefined && liveCacheDims !== embeddingIdentityForRun.dimensions) { + log( + `Embedding dimensions changed (${liveCacheDims}d -> ` + + `${embeddingIdentityForRun.dimensions}d), discarding published cache`, + ); + discardCachedEmbeddings(); + } + if (recovered.rows.length > 0) { + const merged = mergeRecoveredEmbeddings(cachedSnapshot, recovered); + disposeEmbeddingSpill(cachedSnapshot.spill); + adoptCachedEmbeddings(merged); + } + log( + `Recovered ${recovered.rows.length} complete staged embedding chunk(s) ` + + `for ${recovered.embeddingNodeIds.size} node(s); unchanged content can reuse them.`, + ); + } catch (err) { + log( + `Warning: could not recover staged embeddings (${(err as Error).message}); ` + + 'the retry will regenerate missing chunks.', + ); + } finally { + disposeEmbeddingSpill(recovered?.spill); + } + } else { + log('Staged embedding identity or schema changed; its vectors will not be reused.'); + } + } + // ── Load incremental parse cache ────────────────────────────────── // Content-addressed: `--force` reuses parser shards; `useParseCache: false` // stages a new generation under a run-unique parse-rebuild.* dir and publishes @@ -4487,6 +4569,16 @@ async function runFullAnalysisInner( embeddingIdentityForRun = resolveEmbeddingIdentity(); } const embeddingIdentity = embeddingIdentityForRun; + const stagedRecoveryEnabled = useAtomicSwap && isManualCheckpointEnabled(); + if (useAtomicSwap && !stagedRecoveryEnabled) { + log( + 'Manual WAL checkpoints are disabled; new staged work cannot be recovered after interruption. ' + + 'Any previous durable recovery source is retained until publication.', + ); + } + const unsafeRecoveryNodeIds = new Set([...inheritedUnsafeNodeIds, ...restoreFailedNodeIds]); + let activeWindowNodeIds: string[] = []; + let recoveryGenerationDurable = false; // Build a Map from cached embeddings for incremental mode let existingEmbeddings: Map | undefined; if (cachedSnapshot.embeddingNodeIds.size > 0) { @@ -4540,6 +4632,14 @@ async function runFullAnalysisInner( stagedCheckpointEmbeddingCount = embeddings; } const latestMeta = (await loadMeta(metaDir)) ?? existingMeta; + // An in-place write or manual-checkpoint opt-out cannot create a new + // recoverable staged generation. Keep the complete previous receipt: + // updated progress or unsafe nodes would describe different source bytes. + const preservedRecoveryCheckpoint = + !stagedRecoveryEnabled && + resolveEmbeddingRecovery(metaDir, latestMeta?.embeddingCheckpoint) + ? latestMeta?.embeddingCheckpoint + : undefined; // First-ever analyze of this repo: no meta exists on disk yet (the // pre-wipe dirty stamp only fires when one does). Mint the minimum // RepoMeta requires, with `lastCommit: ''` — never `currentCommit` — @@ -4550,6 +4650,11 @@ async function runFullAnalysisInner( lastCommit: '', indexedAt: new Date().toISOString(), }; + const interrupted = mintInterruptedCheckpoint( + embeddingIdentity, + checkpoint, + pendingNodeIds, + ); await saveMeta(metaDir, { ...base, ...(embeddings === undefined || buildPath !== lbugPath @@ -4557,11 +4662,18 @@ async function runFullAnalysisInner( : { stats: { ...base.stats, embeddings } }), // Written by a run that is still IN FLIGHT — see the `kind` doc in // repo-manager.ts. - embeddingCheckpoint: mintInterruptedCheckpoint( - embeddingIdentity, - checkpoint, - pendingNodeIds, - ), + embeddingCheckpoint: preservedRecoveryCheckpoint ?? { + ...interrupted, + ...(stagedRecoveryEnabled + ? { + recovery: { + stagingFile: path.basename(buildPath), + schemaFingerprint: SCHEMA_FINGERPRINT, + unsafeNodeIds: [...unsafeRecoveryNodeIds], + }, + } + : {}), + }, }); }; @@ -4584,7 +4696,21 @@ async function runFullAnalysisInner( { forceReembedNodeIds: pendingEmbeddingNodeIds, onCheckpointWindowStart: async ({ nodeIds, ...checkpoint }) => { + const handoff = stagedRecoveryEnabled && !recoveryGenerationDurable; + if (handoff) { + if (!(await checkpointOnce())) { + throw new Error( + 'Could not checkpoint restored embeddings before recovery handoff.', + ); + } + recoveryGenerationDurable = true; + } + activeWindowNodeIds = nodeIds; + for (const id of nodeIds) unsafeRecoveryNodeIds.add(id); await saveEmbeddingCheckpoint(checkpoint, nodeIds); + // Reclaim the old source only after the new durable generation's + // reference is saved. Later windows retain this same generation. + if (handoff) sweepStagingArtifacts(metaDir, log); }, // ── The mid-run count is a DIAGNOSTIC, not a gate (#2790) ────── // This used to run the count query bare. THIS callback's rejection @@ -4598,7 +4724,12 @@ async function runFullAnalysisInner( // touch stats.embeddings" signal — so the checkpoint still lands, // with whatever count is already on disk left alone. onCheckpoint: async (checkpoint) => { - await checkpointOnce(); + const durable = await checkpointOnce(); + if (stagedRecoveryEnabled && !durable) { + throw new Error('Could not checkpoint the completed embedding window for recovery.'); + } + for (const id of activeWindowNodeIds) unsafeRecoveryNodeIds.delete(id); + activeWindowNodeIds = []; const measured = await measurePersistedEmbeddingCount(executeQuery); if (measured.kind === 'unknown') { log( @@ -5085,6 +5216,7 @@ async function runFullAnalysisInner( // is a crash-safety improvement: a failed swap leaves the previous index // live and the next run recovers via the full-rebuild path. await saveMeta(metaDir, meta); + sweepStagingArtifacts(metaDir, log); // Registry freshness is published only after the graph and its metadata. // A failed close, swap, or metadata save must leave the previous registry @@ -5294,7 +5426,13 @@ async function runFullAnalysisInner( // rethrow below is the surface, and the lock's sweep remains the backstop. if (useAtomicSwap && buildPath !== lbugPath) { try { - await wipeLbugDbFiles(buildPath); + const recovery = resolveEmbeddingRecovery( + metaDir, + (await loadMeta(metaDir))?.embeddingCheckpoint, + ); + // Both paths belong to this locked slot. The validated generation + // basename identifies the same file even through a directory alias. + if (recovery?.stagingFile !== path.basename(buildPath)) await wipeLbugDbFiles(buildPath); } catch { /* swallow — orphan reclamation must never mask the real failure */ } @@ -5306,6 +5444,12 @@ async function runFullAnalysisInner( // IndexLockTimeoutError and other domain failures with `instanceof`. recordLiveIndexMutationRisk(err); } + if (/max(?:imum)?(?: database| db)? size|database size limit|maxDBSize/i.test(String(err))) { + log( + 'The database size limit was reached. Set GITNEXUS_LBUG_MAX_DB_SIZE to a larger ' + + 'byte limit before retrying analyze; retained complete embeddings can be reused.', + ); + } throw err; } } diff --git a/gitnexus/src/server/api.ts b/gitnexus/src/server/api.ts index 86cafb0c1..f4cb39c66 100644 --- a/gitnexus/src/server/api.ts +++ b/gitnexus/src/server/api.ts @@ -12,6 +12,7 @@ import { acquireIndexLock, IndexLockTimeoutError, requireExclusiveIndexLock, + sweepStagingArtifacts, type IndexLockHandle, } from '../storage/index-lock.js'; import { ensurePrivateSharedGraph } from '../core/shared-store-analyze.js'; @@ -2152,11 +2153,21 @@ export const createServer = async (port: number, host: string = '127.0.0.1') => // for the whole embedding write, released in the finally below. let slotLock: IndexLockHandle | undefined; try { - slotLock = await acquireIndexLock(storagePath); + slotLock = await acquireIndexLock(storagePath, { sweep: false }); requireExclusiveIndexLock( slotLock, `Cannot acquire the index lock at ${storagePath}; refusing an unlocked embedding run.`, ); + // This writer cannot recover staged generations. Preserve their + // receipts, including malformed ones, before sweeping or writing. + const recoveryCheckpoint = (await loadMeta(storagePath))?.embeddingCheckpoint; + if (recoveryCheckpoint && Object.hasOwn(recoveryCheckpoint, 'recovery')) { + throw new Error( + 'Cannot generate embeddings: the index checkpoint references staged embeddings. ' + + 'Run `gitnexus analyze` to recover them first.', + ); + } + sweepStagingArtifacts(storagePath); // Writes go to the slot's own graph; a shared-store checkout // reading an immutable commit graph (#3352) takes a private copy. if (!(await ensurePrivateSharedGraph(storagePath, () => {}))) { diff --git a/gitnexus/src/storage/embedding-recovery.ts b/gitnexus/src/storage/embedding-recovery.ts new file mode 100644 index 000000000..54821bcc7 --- /dev/null +++ b/gitnexus/src/storage/embedding-recovery.ts @@ -0,0 +1,167 @@ +/** + * Filesystem-only staged embedding provenance. Keep this independent of native + * and model imports: every index-lock caller needs the retention decision. + */ +import { + closeSync, + constants, + fstatSync, + lstatSync, + openSync, + readFileSync, + realpathSync, +} from 'node:fs'; +import path from 'node:path'; +import type { EmbeddingRecoveryReference } from './repo-meta.js'; +import { INDEX_METADATA_FILE, LEGACY_METADATA_FILE } from './storage-constants.js'; + +const STAGING_FILENAME = + /^lbug\.staging\.[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/; +export const FAMILY_SUFFIXES = [ + '', + '.wal', + '.shadow', + '.wal.checkpoint', + '.lock', + '.checkpoint.intent.lock', + '.checkpoint.apply.lock', +] as const; + +export interface ResolvedEmbeddingRecovery extends EmbeddingRecoveryReference { + dbPath: string; + /** Exact basenames, never a prefix match that can preserve another generation. */ + familyFiles: string[]; +} + +const isRecord = (value: unknown): value is Record => + value !== null && typeof value === 'object' && !Array.isArray(value); +const isNonemptyString = (value: unknown): value is string => + typeof value === 'string' && value.length > 0; +const isNodeIds = (value: unknown): value is string[] => + Array.isArray(value) && value.every(isNonemptyString); +const isCount = (value: unknown): value is number => + typeof value === 'number' && Number.isSafeInteger(value) && value >= 0; + +/** + * Resolve only an explicit interrupted-generation receipt in this canonical + * slot. This validates provenance, not native contents or current identity; + * those checks must pass separately before any rows can be reused. + */ +export const resolveEmbeddingRecovery = ( + lockDir: string, + checkpoint: unknown, +): ResolvedEmbeddingRecovery | undefined => { + if (!isRecord(checkpoint) || !isRecord(checkpoint.recovery)) return undefined; + if (checkpoint.kind !== undefined && checkpoint.kind !== 'interrupted') return undefined; + if ( + !isNonemptyString(checkpoint.at) || + !Number.isFinite(Date.parse(checkpoint.at)) || + !isCount(checkpoint.nodesProcessed) || + !isCount(checkpoint.totalNodes) || + checkpoint.nodesProcessed > checkpoint.totalNodes || + !isCount(checkpoint.chunksProcessed) || + !isNonemptyString(checkpoint.model) || + !isCount(checkpoint.dimensions) || + checkpoint.dimensions === 0 || + !isNonemptyString(checkpoint.provider) || + (checkpoint.pendingNodeIds !== undefined && !isNodeIds(checkpoint.pendingNodeIds)) + ) { + return undefined; + } + const recovery = checkpoint.recovery; + if ( + !isNonemptyString(recovery.stagingFile) || + !STAGING_FILENAME.test(recovery.stagingFile) || + !isNonemptyString(recovery.schemaFingerprint) || + !isNodeIds(recovery.unsafeNodeIds) + ) { + return undefined; + } + const unsafeNodeIds = new Set(recovery.unsafeNodeIds); + if ( + isNodeIds(checkpoint.pendingNodeIds) && + checkpoint.pendingNodeIds.some((nodeId) => !unsafeNodeIds.has(nodeId)) + ) { + return undefined; + } + + try { + const canonicalDir = realpathSync(lockDir); + if (!lstatSync(canonicalDir).isDirectory()) return undefined; + const familyFiles = FAMILY_SUFFIXES.map((suffix) => recovery.stagingFile + suffix); + for (const [index, filename] of familyFiles.entries()) { + try { + // lstat refuses both live and dangling symlinks, without following one + // to a database outside the slot. The base file must exist. + if (!lstatSync(path.join(canonicalDir, filename)).isFile()) return undefined; + } catch (error) { + if (index > 0 && (error as NodeJS.ErrnoException).code === 'ENOENT') continue; + return undefined; + } + } + return { + dbPath: path.join(canonicalDir, recovery.stagingFile), + stagingFile: recovery.stagingFile, + schemaFingerprint: recovery.schemaFingerprint, + unsafeNodeIds: [...unsafeNodeIds], + familyFiles, + }; + } catch { + return undefined; + } +}; + +/** Synchronous mirror of loadMeta's primary-first, absent-only fallback rule. */ +export const readEmbeddingRecovery = (lockDir: string): ResolvedEmbeddingRecovery | undefined => { + let metadataPath = path.join(lockDir, INDEX_METADATA_FILE); + let descriptor: number | undefined; + const noFollow = constants.O_NOFOLLOW ?? 0; + const flags = constants.O_RDONLY | noFollow | (constants.O_NONBLOCK ?? 0); + try { + try { + descriptor = openSync(metadataPath, flags); + } catch (error) { + const code = (error as NodeJS.ErrnoException).code; + if (code !== 'ENOENT' && code !== 'ENOTDIR') return undefined; + // Windows cannot open with O_NOFOLLOW: an open of a dangling symlink + // reports ENOENT, but that existing primary entry must prevent fallback. + try { + lstatSync(metadataPath); + return undefined; + } catch (statError) { + const statCode = (statError as NodeJS.ErrnoException).code; + if (statCode !== 'ENOENT' && statCode !== 'ENOTDIR') return undefined; + } + metadataPath = path.join(lockDir, LEGACY_METADATA_FILE); + descriptor = openSync(metadataPath, flags); + } + const opened = fstatSync(descriptor, { bigint: true }); + const entry = lstatSync(metadataPath, { bigint: true }); + // Check the opened file itself and match the current non-symlink entry. + // This also refuses replacement on platforms without O_NOFOLLOW. + if ( + !opened.isFile() || + !entry.isFile() || + (noFollow === 0 && opened.ino === 0n) || + opened.dev !== entry.dev || + opened.ino !== entry.ino + ) { + return undefined; + } + const meta: unknown = JSON.parse(readFileSync(descriptor, 'utf8')); + if (!isRecord(meta)) return undefined; + // storagePath describes the flat/cache root, including in branch-slot + // metadata. The current locked directory is the generation boundary. + return resolveEmbeddingRecovery(lockDir, meta.embeddingCheckpoint); + } catch { + return undefined; + } finally { + if (descriptor !== undefined) { + try { + closeSync(descriptor); + } catch { + /* best-effort */ + } + } + } +}; diff --git a/gitnexus/src/storage/index-lock.ts b/gitnexus/src/storage/index-lock.ts index 5588127e6..5d3b6cc47 100644 --- a/gitnexus/src/storage/index-lock.ts +++ b/gitnexus/src/storage/index-lock.ts @@ -62,6 +62,7 @@ import path from 'node:path'; import os from 'node:os'; import { randomBytes, randomUUID, createHash } from 'node:crypto'; import { isProcessAlive } from '../utils/process-identity.js'; +import { readEmbeddingRecovery } from './embedding-recovery.js'; const LOCK_FILENAME = 'analyze.lock'; const LOCK_RECORD_VERSION = 1 as const; @@ -433,8 +434,9 @@ const deniedCreateHandle = ( /** * Delete orphaned build/staging artifacts left in the lock directory by a * crashed prior writer. Safe precisely because we hold the exclusive lock: no - * other writer can be creating these here right now, so anything present is a - * crash orphan. Matches this slot's staging files ONLY — never `lbug` itself, + * other writer can be creating these here right now. A checkpoint-referenced + * generation is retained for embedding recovery; all other stages are orphans. + * Matches this slot's staging files ONLY — never `lbug` itself, * never `lbug.wal`/`lbug.shadow` (the LIVE index's own sidecars), and never a * `branches//` sub-slot (which owns its own lock + sweep). Non-recursive. */ @@ -442,6 +444,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo // Matches `lbug.new`, `lbug.new.wal`, `lbug.staging.`, `lbug.staging..wal`, … // Does NOT match `lbug`, `lbug.wal`, `lbug.shadow`. const stagingRe = /^lbug\.(staging\..+|new(\..+)?)$/; + const retained = new Set(readEmbeddingRecovery(lockDir)?.familyFiles ?? []); let removed = 0; let entries: string[]; try { @@ -450,7 +453,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo return; } for (const name of entries) { - if (!stagingRe.test(name)) continue; + if (!stagingRe.test(name) || retained.has(name)) continue; try { unlinkSync(path.join(lockDir, name)); removed++; diff --git a/gitnexus/src/storage/repo-meta.ts b/gitnexus/src/storage/repo-meta.ts index 0c412dd7c..1562ee8c7 100644 --- a/gitnexus/src/storage/repo-meta.ts +++ b/gitnexus/src/storage/repo-meta.ts @@ -43,6 +43,15 @@ export type ContentRetention = 'full' | 'symbol' | 'none'; export type FtsProfile = 'full' | 'symbol-no-file-content' | 'name-only'; export const CONTENT_RETENTION_SCHEMA_VERSION = 1; +/** Exact staged generation whose completed embedding groups can survive a retry. */ +export interface EmbeddingRecoveryReference { + /** A run-minted basename within this metadata file's index slot. */ + stagingFile: string; + schemaFingerprint: string; + /** Active-window and inherited incomplete groups: never reusable until completed. */ + unsafeNodeIds: string[]; +} + /** * Versioned receipt for the analyzer process that produced an index. * @@ -521,6 +530,11 @@ export interface RepoMeta { * subset of their chunks; for `'partial'` they hold none. */ pendingNodeIds?: string[]; + /** + * Interrupted atomic builds only. This does not publish the staged graph + * or advance live embedding statistics; it identifies a recovery source. + */ + recovery?: EmbeddingRecoveryReference; }; /** * Name of the git branch this index represents (#2106). Absent for the diff --git a/gitnexus/test/fixtures/staged-embedding-recovery/seed.mjs b/gitnexus/test/fixtures/staged-embedding-recovery/seed.mjs new file mode 100644 index 000000000..68401a4b8 --- /dev/null +++ b/gitnexus/test/fixtures/staged-embedding-recovery/seed.mjs @@ -0,0 +1,52 @@ +import lbug from '@ladybugdb/core'; +import fs from 'node:fs'; +import { createLbugDatabase } from '../../../src/core/lbug/lbug-config.ts'; + +const [dbPath, mode] = process.argv.slice(2); +const db = createLbugDatabase(lbug, dbPath); +const conn = new lbug.Connection(db); +async function query(cypher) { + const queried = await conn.query(cypher); + for (const result of Array.isArray(queried) ? queried : [queried]) { + await result.getAll(); + await result.close(); + } +} +await query('CREATE NODE TABLE CodeEmbedding (id STRING, nodeId STRING, chunkIndex INT32, startLine INT64, endLine INT64, embedding FLOAT[2], contentHash STRING, PRIMARY KEY(id))'); +async function row(id, nodeId, chunkIndex, hash = 'same', startLine = 1, endLine = 3) { + await query(`CREATE (:CodeEmbedding {id: '${id}', nodeId: '${nodeId}', chunkIndex: ${chunkIndex}, startLine: ${startLine}, endLine: ${endLine}, embedding: [1.0, 2.0], contentHash: ${hash === null ? 'NULL' : `'${hash}'`}})`); +} +await row('complete-0', 'complete', 0); +await row('complete-1', 'complete', 1); +await row('other', 'other', 0); +await query('CHECKPOINT'); +if (mode === 'interrupted-checkpoint') { + await row('checkpoint-only', 'checkpoint-only', 0); + const main = fs.readFileSync(dbPath); + const wal = fs.readFileSync(`${dbPath}.wal`); + await conn.close(); + await db.close(); + // Restore the pre-close bytes: writable close has already checkpointed them. + fs.writeFileSync(dbPath, main); + fs.writeFileSync(`${dbPath}.wal.checkpoint`, wal); + fs.writeFileSync(`${dbPath}.wal`, ''); + fs.writeFileSync(`${dbPath}.shadow`, ''); + fs.writeFileSync(`${dbPath}.checkpoint.intent.lock`, ''); + fs.writeFileSync(`${dbPath}.checkpoint.apply.lock`, ''); +} else if (mode === 'hard-kill') { + await row('unsafe', 'unsafe-prefix', 0); + process.kill(process.pid, 'SIGKILL'); +} else { + await row('gap', 'gap', 1); + await row('duplicate-0', 'duplicate', 0); + await row('duplicate-1', 'duplicate', 0); + await row('mixed-0', 'mixed', 0, 'old'); + await row('mixed-1', 'mixed', 1, 'new'); + await row('missing-hash-0', 'missing-hash', 0); + await row('missing-hash-1', 'missing-hash', 1, null); + await row('bad-line', 'bad-line', 0, 'same', -1, 3); + await row('nan-0', 'nan', 0); + await query("CREATE (:CodeEmbedding {id: 'nan-1', nodeId: 'nan', chunkIndex: 1, startLine: 1, endLine: 3, embedding: [CAST('NaN', 'FLOAT'), 2.0], contentHash: 'same'})"); + await conn.close(); + await db.close(); +} diff --git a/gitnexus/test/integration/analyze-staged-embedding-recovery.test.ts b/gitnexus/test/integration/analyze-staged-embedding-recovery.test.ts new file mode 100644 index 000000000..dbb5b893e --- /dev/null +++ b/gitnexus/test/integration/analyze-staged-embedding-recovery.test.ts @@ -0,0 +1,449 @@ +/** + * Exercise the CLI -> native staged DB -> durable checkpoint -> SIGKILL -> + * isolated recovery -> graph rebuild -> publication chain. The endpoint is + * local and deterministic; request text is the billing/reuse oracle. + */ +import { spawn, spawnSync, type ChildProcess } from 'node:child_process'; +import fs from 'node:fs'; +import http from 'node:http'; +import os from 'node:os'; +import path from 'node:path'; +import { fileURLToPath, pathToFileURL } from 'node:url'; +import { afterAll, beforeAll, describe, expect, it } from 'vitest'; +import { CLI_SPAWN_PREFIX, tsxLoaderUrl } from '../helpers/cli-entry.js'; + +const DIMS = 8; +const NODE_COUNT = 5_128; +const DEADLINE = process.env.CI ? 180_000 : 120_000; +const packageRoot = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..', '..'); +const adapterUrl = pathToFileURL(path.join(packageRoot, 'src/core/lbug/lbug-adapter.ts')).href; + +interface CheckpointMeta { + repoPath: string; + stats?: { embeddings?: number }; + embeddingCheckpoint?: { + nodesProcessed: number; + chunksProcessed: number; + pendingNodeIds?: string[]; + recovery?: { stagingFile: string; schemaFingerprint: string; unsafeNodeIds: string[] }; + }; +} + +interface Run { + child: ChildProcess; + output: () => string; + done: Promise<{ code: number | null; signal: NodeJS.Signals | null; output: string }>; +} + +let root: string; +let stoppedRepo: string; +let server: http.Server; +let endpoint: string; +let submitted: string[] = []; +let completed: string[] = []; +let durableTexts: string[] = []; +const activeWindowCompleted: string[] = []; +let held: string[] = []; +let stopAfterCheckpoint = false; +let currentRepo: string; +let gateResolve: (() => void) | undefined; +const running = new Set(); + +function readMeta(repo: string): CheckpointMeta { + return JSON.parse(fs.readFileSync(path.join(repo, '.gitnexus', 'gitnexus.json'), 'utf8')); +} + +function cliEnv(repo: string): NodeJS.ProcessEnv { + const env = { ...process.env }; + for (const key of Object.keys(env)) { + if (key.startsWith('GITNEXUS_EMBEDDING_') || key.startsWith('GITNEXUS_STORAGE_')) + delete env[key]; + } + return { + ...env, + GITNEXUS_HOME: path.join(root, `home-${path.basename(repo)}`), + GITNEXUS_SHARED_STORE: 'off', + GITNEXUS_LBUG_EXTENSION_INSTALL: 'never', + GITNEXUS_LBUG_BUFFER_POOL_SIZE: String(256 * 1024 * 1024), + GITNEXUS_EMBEDDING_URL: endpoint, + GITNEXUS_EMBEDDING_MODEL: 'staged-recovery-fixture', + GITNEXUS_EMBEDDING_DIMS: String(DIMS), + GITNEXUS_EMBEDDING_BATCH_SIZE: '5000', + GITNEXUS_EMBEDDING_SUB_BATCH_SIZE: '64', + GITNEXUS_EMBEDDING_MAX_ATTEMPTS: '1', + GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT: '0', + GITNEXUS_MEMORY: 'off', + // SIGKILL skips process exit hooks; keep orphaned cache/export spills in + // this suite's owned directory so afterAll can remove them as well. + TMPDIR: path.join(root, 'tmp'), + NODE_OPTIONS: `${process.env.NODE_OPTIONS || ''} --max-old-space-size=2048`.trim(), + CI: '1', + }; +} + +function runAnalyze(repo: string, flags: string[] = [], onOutput?: (output: string) => void): Run { + let output = ''; + const child = spawn( + process.execPath, + [ + ...CLI_SPAWN_PREFIX, + 'analyze', + repo, + '--no-share', + '--skip-skills', + '--skip-fts', + '--workers', + '1', + ...flags, + ], + { cwd: repo, env: cliEnv(repo), stdio: ['ignore', 'pipe', 'pipe'] }, + ); + running.add(child); + const timeout = setTimeout(() => child.kill('SIGKILL'), DEADLINE); + const collect = (chunk: Buffer) => { + output += chunk.toString(); + onOutput?.(output); + }; + child.stdout?.on('data', collect); + child.stderr?.on('data', collect); + const done = new Promise<{ code: number | null; signal: NodeJS.Signals | null; output: string }>( + (resolve, reject) => { + child.once('error', reject); + child.once('close', (code, signal) => { + clearTimeout(timeout); + running.delete(child); + resolve({ code, signal, output }); + }); + }, + ); + return { child, done, output: () => output }; +} + +async function successfulAnalyze(repo: string, flags: string[] = []): Promise { + submitted = []; + stopAfterCheckpoint = false; + currentRepo = repo; + const result = await runAnalyze(repo, flags).done; + expect(result.output, `CLI exited ${result.code}, signal ${result.signal}`).not.toContain( + 'SIGABRT', + ); + expect(result.code, result.output).toBe(0); + return [...submitted]; +} + +function cloneStoppedRepo(name: string): string { + const repo = path.join(root, name); + fs.cpSync(stoppedRepo, repo, { recursive: true }); + for (const filename of ['gitnexus.json', 'meta.json']) { + const target = path.join(repo, '.gitnexus', filename); + const meta = JSON.parse(fs.readFileSync(target, 'utf8')); + meta.repoPath = repo; + meta.storagePath = path.join(repo, '.gitnexus'); + fs.writeFileSync(target, JSON.stringify(meta)); + } + return repo; +} + +function readPublishedRows( + repo: string, +): Array<{ nodeId: string; chunkIndex: number; embedding: number[] }> { + // Keep native handles out of the vitest fork, and wait for clean teardown + // before accepting the receipt. Every read opens only a published DB. + const receiptPath = path.join(root, `rows-${path.basename(repo)}.json`); + const script = ` + const adapter = await import(${JSON.stringify(adapterUrl)}); + const fs = await import('node:fs'); + await adapter.initLbug(${JSON.stringify(path.join(repo, '.gitnexus', 'lbug'))}); + try { + const rows = await adapter.executeQuery('MATCH (e:CodeEmbedding) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.embedding AS embedding'); + fs.writeFileSync(${JSON.stringify(receiptPath)}, JSON.stringify(rows)); + console.log('ROWS_RECEIPT:' + rows.length); + } finally { await adapter.closeLbug(); } + `; + const result = spawnSync( + process.execPath, + ['--import', tsxLoaderUrl(), '--input-type=module', '-e', script], + { + cwd: packageRoot, + env: cliEnv(repo), + encoding: 'utf8', + timeout: 30_000, + }, + ); + try { + const diagnostic = + `${result.error ?? ''} ${result.signal ?? ''}\n${result.stderr}\n${result.stdout}`.slice( + 0, + 4096, + ); + expect(result.status, diagnostic).toBe(0); + expect(result.stdout).toMatch(/ROWS_RECEIPT:\d+/); + return JSON.parse(fs.readFileSync(receiptPath, 'utf8')); + } finally { + fs.rmSync(receiptPath, { force: true }); + } +} + +function expectPublishedComplete(repo: string, expectedNodes: number): void { + const meta = readMeta(repo); + expect(meta.embeddingCheckpoint).toBeUndefined(); + const rows = readPublishedRows(repo); + expect(new Set(rows.map((row) => row.nodeId)).size).toBe(expectedNodes); + expect(meta.stats?.embeddings).toBe(rows.length); + const byNode = new Map(); + for (const row of rows) { + expect(row.embedding).toHaveLength(DIMS); + expect(row.embedding.every(Number.isFinite)).toBe(true); + const indices = byNode.get(row.nodeId) ?? []; + indices.push(row.chunkIndex); + byNode.set(row.nodeId, indices); + } + for (const indices of byNode.values()) { + expect(indices.sort((a, b) => a - b)).toEqual( + Array.from({ length: indices.length }, (_, index) => index), + ); + } + expect( + fs.readdirSync(path.join(repo, '.gitnexus')).filter((name) => name.startsWith('lbug.staging.')), + ).toEqual([]); +} + +beforeAll(async () => { + if (process.platform === 'win32') return; + root = fs.mkdtempSync(path.join(os.tmpdir(), 'gn-staged-recovery-e2e-')); + fs.mkdirSync(path.join(root, 'tmp')); + stoppedRepo = path.join(root, 'interrupted'); + currentRepo = stoppedRepo; + fs.mkdirSync(stoppedRepo); + const shortFunctions = Array.from( + { length: NODE_COUNT - 1 }, + (_, index) => `export function recoverable${index}() { return ${index}; }`, + ); + // Multi-chunk nodes must be reused as a complete group, including their tail. + const longFunction = `export function longRecoverable() {\n${Array.from({ length: 80 }, (_, index) => ` // retained chunk marker ${index} ${'x'.repeat(80)}`).join('\n')}\n return 42;\n}`; + fs.writeFileSync( + path.join(stoppedRepo, 'functions.ts'), + `${longFunction}\n${shortFunctions.join('\n')}\n`, + ); + const gitEnv = { + ...process.env, + GIT_AUTHOR_NAME: 'test', + GIT_AUTHOR_EMAIL: 'test@test', + GIT_COMMITTER_NAME: 'test', + GIT_COMMITTER_EMAIL: 'test@test', + }; + for (const args of [['init'], ['add', 'functions.ts'], ['commit', '-m', 'recovery fixture']]) { + const result = spawnSync('git', args, { cwd: stoppedRepo, env: gitEnv, encoding: 'utf8' }); + expect(result.status, result.stderr).toBe(0); + } + server = http.createServer((request, response) => { + let body = ''; + request.on('data', (chunk) => { + body += chunk; + }); + request.on('end', () => { + const { input } = JSON.parse(body) as { input: string[] }; + submitted.push(...input); + if ( + stopAfterCheckpoint && + (readMeta(currentRepo).embeddingCheckpoint?.nodesProcessed ?? 0) >= 5000 + ) { + if (activeWindowCompleted.length > 0) { + held = [...input]; + gateResolve?.(); + // One sub-batch has already inserted rows inside the unsafe window. + // Awaiting the next real request gates a crash before it completes. + return; + } + durableTexts = [...completed]; + activeWindowCompleted.push(...input); + } + completed.push(...input); + response.writeHead(200, { 'content-type': 'application/json' }); + response.end( + JSON.stringify({ + data: input.map((text, index) => ({ + index, + embedding: Array.from( + { length: DIMS }, + (_, dimension) => ((text.length + dimension * 17) % 101) / 101, + ), + })), + }), + ); + }); + }); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + const address = server.address() as { port: number }; + endpoint = `http://127.0.0.1:${address.port}/v1`; + await successfulAnalyze(stoppedRepo, ['--index-only']); + expect(readMeta(stoppedRepo).stats?.embeddings ?? 0).toBe(0); + completed = []; + held = []; + stopAfterCheckpoint = true; + const gate = new Promise((resolve) => { + gateResolve = resolve; + }); + const first = runAnalyze(stoppedRepo, ['--force', '--embeddings']); + await Promise.race([ + gate, + first.done.then((result) => { + throw new Error(`CLI exited before recovery gate: ${result.output}`); + }), + ]); + first.child.kill('SIGKILL'); + expect((await first.done).signal).toBe('SIGKILL'); + const interrupted = readMeta(stoppedRepo); + expect(interrupted.embeddingCheckpoint?.nodesProcessed).toBe(5000); + expect(interrupted.embeddingCheckpoint?.recovery?.stagingFile).toMatch(/^lbug\.staging\./); + expect(interrupted.stats?.embeddings ?? 0).toBe(0); + expect(completed.length).toBeGreaterThanOrEqual(5000); + expect(completed.filter((text) => text.includes('longRecoverable')).length).toBeGreaterThan(1); + expect(activeWindowCompleted).toHaveLength(64); + expect(held.length).toBeGreaterThan(0); + gateResolve = undefined; + stopAfterCheckpoint = false; +}, DEADLINE * 2); + +afterAll(async () => { + for (const child of running) child.kill('SIGKILL'); + await Promise.all( + [...running].map( + (child) => new Promise((resolve) => child.once('close', () => resolve())), + ), + ); + if (server) { + server.closeAllConnections(); + await new Promise((resolve) => server.close(() => resolve())); + } + if (root) fs.rmSync(root, { recursive: true, force: true }); +}); + +// Atomic publication is POSIX-specific; Windows uses the in-place path. +describe + .skipIf(process.platform === 'win32') + .sequential('interrupted staged embedding recovery (real CLI and native DB)', () => { + it.each([ + ['plain', []], + ['forced', ['--force', '--embeddings']], + ] as const)( + '%s retry bills only unfinished groups and publishes an honest count', + async (name, flags) => { + const repo = cloneStoppedRepo(name); + const texts = await successfulAnalyze(repo, [...flags]); + const durable = new Set(durableTexts); + expect(texts.filter((text) => durable.has(text))).toEqual([]); + for (const text of held) expect(texts).toContain(text); + for (const text of activeWindowCompleted) expect(texts).toContain(text); + expect(texts.length).toBe(128); + expectPublishedComplete(repo, NODE_COUNT); + }, + DEADLINE, + ); + + it( + 'manual checkpoint opt-out completes a staged retry without resubmitting durable chunks', + async () => { + const repo = cloneStoppedRepo('manual-checkpoint-opt-out'); + const previous = process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT; + process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT = '0'; + try { + const texts = await successfulAnalyze(repo, ['--force', '--embeddings']); + expect(texts.filter((text) => new Set(durableTexts).has(text))).toEqual([]); + expect(texts.length).toBe(128); + expectPublishedComplete(repo, NODE_COUNT); + } finally { + if (previous === undefined) delete process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT; + else process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT = previous; + } + }, + DEADLINE, + ); + + it( + 'survives another crash after harvesting but before the replacement is durable', + async () => { + const repo = cloneStoppedRepo('crash-again'); + const source = readMeta(repo).embeddingCheckpoint?.recovery?.stagingFile; + expect(source).toBeDefined(); + if (!source) throw new Error('fixture has no retained generation'); + submitted = []; + currentRepo = repo; + let killed = false; + const retry = runAnalyze(repo, [], (output) => { + if (!killed && /Recovered \d+ complete staged embedding chunk/.test(output)) { + killed = true; + retry.child.kill('SIGKILL'); + } + }); + const result = await retry.done; + expect(killed, result.output).toBe(true); + expect(result.signal).toBe('SIGKILL'); + expect(submitted).toEqual([]); + expect(readMeta(repo).embeddingCheckpoint?.recovery?.stagingFile).toBe(source); + expect(fs.existsSync(path.join(repo, '.gitnexus', source))).toBe(true); + const finalTexts = await successfulAnalyze(repo); + expect(finalTexts.filter((text) => new Set(durableTexts).has(text))).toEqual([]); + expect(finalTexts.length).toBe(128); + expectPublishedComplete(repo, NODE_COUNT); + }, + DEADLINE * 2, + ); + + it( + 'regenerates changed content and removes deleted nodes while reusing the other complete groups', + async () => { + const repo = cloneStoppedRepo('changed'); + const source = path.join(repo, 'functions.ts'); + const content = fs + .readFileSync(source, 'utf8') + .replace( + 'export function recoverable5() { return 5; }', + 'export function recoverable5() { return 999999; }', + ) + .replace( + 'export function recoverable6() { return 6; }', + '// deleted function retains line offsets', + ); + fs.writeFileSync(source, content); + const texts = await successfulAnalyze(repo); + expect(texts.some((text) => text.includes('recoverable5') && text.includes('999999'))).toBe( + true, + ); + expect(texts.some((text) => text.includes('recoverable6'))).toBe(false); + expect(texts.filter((text) => new Set(durableTexts).has(text))).toEqual([]); + expect(texts.length).toBe(129); + expectPublishedComplete(repo, NODE_COUNT - 1); + }, + DEADLINE, + ); + + it( + 'a forced retry with a different model does not import the staged cache', + async () => { + const repo = cloneStoppedRepo('different-model'); + const texts = await successfulAnalyze(repo, [ + '--force', + '--embeddings', + '--embedding-model', + 'different-model', + ]); + for (const text of durableTexts) expect(texts).toContain(text); + expect(texts.length).toBeGreaterThanOrEqual(NODE_COUNT); + expectPublishedComplete(repo, NODE_COUNT); + }, + DEADLINE, + ); + + it( + 'explicit drop abandons staged vectors without contacting the provider', + async () => { + const repo = cloneStoppedRepo('drop'); + expect(await successfulAnalyze(repo, ['--force', '--drop-embeddings'])).toEqual([]); + expect(readMeta(repo).embeddingCheckpoint).toBeUndefined(); + expect(readMeta(repo).stats?.embeddings ?? 0).toBe(0); + expect(readPublishedRows(repo)).toEqual([]); + }, + DEADLINE, + ); + }); diff --git a/gitnexus/test/integration/staged-embedding-recovery.test.ts b/gitnexus/test/integration/staged-embedding-recovery.test.ts new file mode 100644 index 000000000..0a52c67ab --- /dev/null +++ b/gitnexus/test/integration/staged-embedding-recovery.test.ts @@ -0,0 +1,208 @@ +import { spawnSync } from 'node:child_process'; +import { randomUUID } from 'node:crypto'; +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { afterEach, describe, expect, it } from 'vitest'; +import { + disposeEmbeddingSpill, + materializeCachedEmbeddings, + type CachedEmbeddingsSnapshot, +} from '../../src/core/embeddings/embedding-restore-spill.js'; +import { recoverStagedEmbeddings } from '../../src/core/embeddings/staged-embedding-recovery.js'; + +describe('isolated native staged embedding recovery', () => { + let tmp: string | undefined; + let recovered: CachedEmbeddingsSnapshot | undefined; + afterEach(() => { + disposeEmbeddingSpill(recovered?.spill); + recovered = undefined; + if (tmp) fs.rmSync(tmp, { recursive: true, force: true }); + tmp = undefined; + }); + function stagePath() { + tmp ??= fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-native-')); + return path.join(tmp, `lbug.staging.${randomUUID()}`); + } + function seed(dbPath: string, mode = 'clean') { + const result = spawnSync( + process.execPath, + [ + '--import', + 'tsx', + fileURLToPath(new URL('../fixtures/staged-embedding-recovery/seed.mjs', import.meta.url)), + dbPath, + mode, + ], + { + encoding: 'utf8', + timeout: 20_000, + env: { ...process.env, GITNEXUS_LBUG_BUFFER_POOL_SIZE: String(128 * 1024 * 1024) }, + }, + ); + expect(result.error, result.stderr).toBeUndefined(); + if (mode === 'hard-kill') expect(result.signal, result.stderr).toBe('SIGKILL'); + else expect(result.status, result.stderr).toBe(0); + } + + function snapshotSourceFamily(dbPath: string) { + const basename = path.basename(dbPath); + return Object.fromEntries( + fs + .readdirSync(path.dirname(dbPath)) + .filter((name) => name === basename || name.startsWith(`${basename}.`)) + .sort() + .map((name) => [name, fs.readFileSync(path.join(path.dirname(dbPath), name))]), + ); + } + + it('streams complete same-hash groups and rejects malformed whole nodes', async () => { + const dbPath = stagePath(); + seed(dbPath); + recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 2 }); + expect([...recovered.embeddingNodeIds].sort()).toEqual(['complete', 'other']); + expect(recovered.rows).toHaveLength(3); + expect(recovered.embeddings).toEqual([]); + expect(materializeCachedEmbeddings(recovered, recovered.rows)).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + nodeId: 'complete', + chunkIndex: 0, + embedding: [1, 2], + contentHash: 'same', + }), + expect.objectContaining({ + nodeId: 'complete', + chunkIndex: 1, + embedding: [1, 2], + contentHash: 'same', + }), + ]), + ); + expect(fs.existsSync(dbPath)).toBe(true); + }, 30_000); + + it('replays a hard-killed native writer strictly and excludes the incomplete active window', async () => { + const dbPath = stagePath(); + seed(dbPath, 'hard-kill'); + const sourceBefore = snapshotSourceFamily(dbPath); + recovered = await recoverStagedEmbeddings(dbPath, { + dimensions: 2, + excludedNodeIds: ['unsafe-prefix', 'other'], + }); + expect([...recovered.embeddingNodeIds]).toEqual(['complete']); + expect(recovered.rows).toHaveLength(2); + expect( + materializeCachedEmbeddings(recovered, recovered.rows) + .map((row) => row.chunkIndex) + .sort(), + ).toEqual([0, 1]); + expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore); + }, 30_000); + + it('rejects vectors from a different dimension instead of coercing them', async () => { + const dbPath = stagePath(); + seed(dbPath); + recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 3 }); + expect(recovered.rows).toEqual([]); + }, 30_000); + + it('strictly replays an interrupted checkpoint copy with both checkpoint locks retained', async (ctx) => { + const version = JSON.parse( + fs.readFileSync( + new URL('../../node_modules/@ladybugdb/core/package.json', import.meta.url), + 'utf8', + ), + ).version as string; + const [major, minor] = version.split('.').map(Number); + // Older pins do not support the deterministic interrupted-checkpoint plant. + if (Number.isFinite(major) && Number.isFinite(minor) && major === 0 && minor < 19) ctx.skip(); + const dbPath = stagePath(); + seed(dbPath, 'interrupted-checkpoint'); + const sourceBefore = snapshotSourceFamily(dbPath); + expect(fs.statSync(`${dbPath}.wal.checkpoint`).size).toBeGreaterThan(0); + expect(fs.existsSync(`${dbPath}.checkpoint.intent.lock`)).toBe(true); + expect(fs.existsSync(`${dbPath}.checkpoint.apply.lock`)).toBe(true); + + recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 2 }); + + expect([...recovered.embeddingNodeIds].sort()).toEqual([ + 'checkpoint-only', + 'complete', + 'other', + ]); + expect(recovered.rows).toHaveLength(4); + expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore); + }, 30_000); + + it('contains a malformed native source in a subprocess and preserves it', async () => { + const dbPath = stagePath(); + fs.writeFileSync(dbPath, 'not a ladybug database'); + await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow( + /extraction failed/, + ); + expect(fs.readFileSync(dbPath, 'utf8')).toBe('not a ladybug database'); + }, 30_000); + + it('rejects a malformed WAL without deleting or quarantining it to reopen the source', async () => { + const dbPath = stagePath(); + seed(dbPath, 'hard-kill'); + const walPath = `${dbPath}.wal`; + fs.writeFileSync(walPath, Buffer.alloc(128, 0xff)); + const sourceBefore = snapshotSourceFamily(dbPath); + await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow( + /extraction failed/, + ); + expect(fs.existsSync(walPath)).toBe(true); + expect(fs.readdirSync(path.dirname(dbPath)).some((name) => /bad|quarantine/i.test(name))).toBe( + false, + ); + expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore); + }, 30_000); + + it('does not create a missing source and refuses a symlink', async () => { + const dbPath = stagePath(); + await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow(/ENOENT/); + expect(fs.existsSync(dbPath)).toBe(false); + const realPath = path.join(path.dirname(dbPath), 'real'); + fs.writeFileSync(realPath, 'fixture'); + fs.symlinkSync(realPath, dbPath); + await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow( + /regular file/, + ); + }); + + it('can kill a timed out native subprocess without aborting analyze', async () => { + const dbPath = stagePath(); + seed(dbPath); + const sourceBefore = snapshotSourceFamily(dbPath); + await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2, timeoutMs: 1 })).rejects.toThrow( + /timeout/, + ); + expect(fs.existsSync(dbPath)).toBe(true); + expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore); + }, 30_000); + + it('loads the source child when analyze runs in another repository directory', () => { + const dbPath = stagePath(); + seed(dbPath); + const result = spawnSync( + process.execPath, + [ + '--import', + import.meta.resolve('tsx'), + '--input-type=module', + '-e', + 'const { recoverStagedEmbeddings } = await import(process.argv[1]); const { disposeEmbeddingSpill } = await import(process.argv[2]); const recovered = await recoverStagedEmbeddings(process.argv[3], {dimensions: 2}); process.stdout.write(String(recovered.rows.length)); disposeEmbeddingSpill(recovered.spill);', + new URL('../../src/core/embeddings/staged-embedding-recovery.ts', import.meta.url).href, + new URL('../../src/core/embeddings/embedding-restore-spill.ts', import.meta.url).href, + dbPath, + ], + { cwd: path.dirname(dbPath), encoding: 'utf8', timeout: 20_000 }, + ); + expect(result.error, result.stderr).toBeUndefined(); + expect(result.status, result.stderr).toBe(0); + expect(result.stdout).toBe('3'); + }, 30_000); +}); diff --git a/gitnexus/test/unit/api-fts-mode.test.ts b/gitnexus/test/unit/api-fts-mode.test.ts index 632bd37d3..ab8eaf401 100644 --- a/gitnexus/test/unit/api-fts-mode.test.ts +++ b/gitnexus/test/unit/api-fts-mode.test.ts @@ -7,7 +7,12 @@ import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } const mocks = vi.hoisted(() => ({ loadMeta: vi.fn(), + saveMeta: vi.fn(), listRegisteredRepos: vi.fn(), + acquireIndexLock: vi.fn(), + releaseIndexLock: vi.fn(), + ensurePrivateSharedGraph: vi.fn(), + runEmbeddingPipeline: vi.fn(), withLbugDb: vi.fn(), search: vi.fn(), updateJob: vi.fn(), @@ -16,8 +21,19 @@ const mocks = vi.hoisted(() => ({ vi.mock('../../src/storage/repo-manager.js', async (importOriginal) => ({ ...(await importOriginal()), loadMeta: mocks.loadMeta, + saveMeta: mocks.saveMeta, listRegisteredRepos: mocks.listRegisteredRepos, })); +vi.mock('../../src/storage/index-lock.js', async (importOriginal) => ({ + ...(await importOriginal()), + acquireIndexLock: mocks.acquireIndexLock, +})); +vi.mock('../../src/core/shared-store-analyze.js', () => ({ + ensurePrivateSharedGraph: mocks.ensurePrivateSharedGraph, +})); +vi.mock('../../src/core/embeddings/embedding-pipeline.js', () => ({ + runEmbeddingPipeline: mocks.runEmbeddingPipeline, +})); vi.mock('../../src/storage/storage-resolver.js', async (importOriginal) => ({ ...(await importOriginal()), requireRegisteredStoragePath: vi.fn(async (entry: { storagePath: string }) => entry.storagePath), @@ -116,6 +132,19 @@ afterAll(() => { beforeEach(() => { vi.clearAllMocks(); + mocks.acquireIndexLock.mockResolvedValue({ + release: mocks.releaseIndexLock, + record: { + v: 1, + pid: process.pid, + hostname: 'test-host', + startTime: null, + token: 'test-lock', + invocationId: 'test-run', + acquiredAt: '', + }, + }); + mocks.ensurePrivateSharedGraph.mockResolvedValue(true); mocks.listRegisteredRepos.mockResolvedValue([entry]); mocks.withLbugDb.mockImplementation(async (_path, callback) => callback()); mocks.search.mockImplementation(async (_query, _limit, _exec, reason) => ({ @@ -124,6 +153,145 @@ beforeEach(() => { })); }); +describe('POST /api/embed staged recovery preflight', () => { + afterEach(() => { + vi.unstubAllEnvs(); + fs.rmSync(entry.storagePath, { recursive: true, force: true }); + fs.mkdirSync(entry.storagePath, { recursive: true }); + }); + + async function useRealIndexLock() { + vi.stubEnv('GITNEXUS_INDEX_LOCK_BACKEND', 'file'); + const actual = await vi.importActual( + '../../src/storage/index-lock.js', + ); + mocks.acquireIndexLock.mockImplementation( + async (...args: Parameters) => { + const lock = await actual.acquireIndexLock(...args); + return { + ...lock, + release: () => { + lock.release(); + mocks.releaseIndexLock(); + }, + }; + }, + ); + } + + it.each([ + { + name: 'valid staged receipt', + recovery: { + stagingFile: 'lbug.staging.12345678-1234-4123-8123-123456789abc', + schemaFingerprint: 'test-schema', + unsafeNodeIds: ['n2', 'inherited-window-node'], + }, + }, + { name: 'null receipt', recovery: null }, + { name: 'malformed receipt', recovery: { stagingFile: 'invalid' } }, + { name: 'false receipt', recovery: false }, + ])('preserves a $name and releases the job locks', async ({ recovery }) => { + await useRealIndexLock(); + const lockPath = path.join(entry.storagePath, 'analyze.lock'); + const lbugPath = path.join(entry.storagePath, 'lbug'); + const metaPath = path.join(entry.storagePath, 'gitnexus.json'); + const sourcePath = path.join( + entry.storagePath, + 'lbug.staging.12345678-1234-4123-8123-123456789abc', + ); + fs.writeFileSync(lbugPath, 'published graph'); + fs.writeFileSync(sourcePath, 'completed paid vectors'); + fs.writeFileSync(`${sourcePath}.wal`, 'unfinished window'); + const metadataBytes = JSON.stringify({ + repoPath: entry.path, + lastCommit: 'abc123', + indexedAt: '2026-01-01T00:00:00.000Z', + stats: { embeddings: 7 }, + embeddingCheckpoint: { + at: '2026-01-01T00:00:00.000Z', + nodesProcessed: 1, + totalNodes: 2, + chunksProcessed: 1, + model: 'test-model', + dimensions: 768, + provider: 'local', + kind: 'interrupted', + pendingNodeIds: ['n2'], + recovery, + }, + }); + fs.writeFileSync(metaPath, metadataBytes); + mocks.loadMeta.mockImplementation(async () => { + expect(fs.existsSync(lockPath)).toBe(true); + return JSON.parse(fs.readFileSync(metaPath, 'utf8')); + }); + mocks.withLbugDb.mockResolvedValue(undefined); + + await invoke('/api/embed'); + await vi.waitFor(() => expect(mocks.releaseIndexLock).toHaveBeenCalledTimes(1)); + + expect(mocks.updateJob).toHaveBeenCalledWith( + 'embed-job', + expect.objectContaining({ + status: 'failed', + error: expect.stringMatching( + /staged embeddings.*Run `gitnexus analyze` to recover them first/, + ), + }), + ); + expect(mocks.acquireIndexLock).toHaveBeenCalledWith(entry.storagePath, { sweep: false }); + expect(mocks.loadMeta.mock.invocationCallOrder[0]).toBeGreaterThan( + mocks.acquireIndexLock.mock.invocationCallOrder[0]!, + ); + expect(mocks.ensurePrivateSharedGraph).not.toHaveBeenCalled(); + expect(mocks.withLbugDb).not.toHaveBeenCalled(); + expect(mocks.runEmbeddingPipeline).not.toHaveBeenCalled(); + expect(mocks.saveMeta).not.toHaveBeenCalled(); + expect(fs.existsSync(lockPath)).toBe(false); + expect(fs.readFileSync(metaPath, 'utf8')).toBe(metadataBytes); + expect(fs.readFileSync(lbugPath, 'utf8')).toBe('published graph'); + expect(fs.readFileSync(sourcePath, 'utf8')).toBe('completed paid vectors'); + expect(fs.readFileSync(`${sourcePath}.wal`, 'utf8')).toBe('unfinished window'); + + // A second accepted job proves the in-memory repo lock was also released. + await invoke('/api/embed'); + await vi.waitFor(() => expect(mocks.releaseIndexLock).toHaveBeenCalledTimes(2)); + expect(fs.existsSync(lockPath)).toBe(false); + }); + + it('sweeps orphaned staging files before a writable job without a recovery receipt', async () => { + await useRealIndexLock(); + const metaPath = path.join(entry.storagePath, 'gitnexus.json'); + fs.writeFileSync(metaPath, JSON.stringify({ repoPath: entry.path })); + const sourcePath = path.join( + entry.storagePath, + 'lbug.staging.12345678-1234-4123-8123-123456789abc', + ); + fs.writeFileSync(sourcePath, 'orphaned database'); + fs.writeFileSync(`${sourcePath}.wal`, 'orphaned WAL'); + mocks.loadMeta.mockImplementation(async () => JSON.parse(fs.readFileSync(metaPath, 'utf8'))); + mocks.ensurePrivateSharedGraph.mockImplementation(async () => { + expect(fs.existsSync(path.join(entry.storagePath, 'analyze.lock'))).toBe(true); + expect(fs.existsSync(sourcePath)).toBe(false); + expect(fs.existsSync(`${sourcePath}.wal`)).toBe(false); + return true; + }); + mocks.withLbugDb.mockResolvedValue(undefined); + + await invoke('/api/embed'); + await vi.waitFor(() => expect(mocks.releaseIndexLock).toHaveBeenCalledTimes(1)); + + expect(mocks.updateJob).toHaveBeenCalledWith( + 'embed-job', + expect.objectContaining({ status: 'complete' }), + ); + expect(mocks.ensurePrivateSharedGraph).toHaveBeenCalledTimes(1); + expect(mocks.withLbugDb).toHaveBeenCalledTimes(1); + expect(fs.existsSync(path.join(entry.storagePath, 'analyze.lock'))).toBe(false); + }); +}); + async function invoke(route: string, query: Record = {}) { const layer = app.router.stack.find((item: any) => item.route?.path === route); expect(layer, route).toBeDefined(); @@ -257,7 +425,9 @@ describe('serve uses one metadata-derived FTS mode on every DB-open path', () => expect.any(Function), skip ? { skipFts: true } : {}, ); - expect(mocks.loadMeta).toHaveBeenCalledExactlyOnceWith(entry.storagePath); + expect(mocks.loadMeta).toHaveBeenCalledTimes(2); + expect(mocks.loadMeta).toHaveBeenNthCalledWith(1, entry.storagePath); + expect(mocks.loadMeta).toHaveBeenNthCalledWith(2, entry.storagePath); }, ); }); diff --git a/gitnexus/test/unit/embedding-recovery-race.test.ts b/gitnexus/test/unit/embedding-recovery-race.test.ts new file mode 100644 index 000000000..145f860ec --- /dev/null +++ b/gitnexus/test/unit/embedding-recovery-race.test.ts @@ -0,0 +1,199 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import * as fs from 'node:fs'; +import { execFileSync } from 'node:child_process'; +import os from 'node:os'; +import path from 'node:path'; +import { readEmbeddingRecovery } from '../../src/storage/embedding-recovery.js'; + +vi.mock('node:fs', async (importOriginal) => { + const actual = await importOriginal(); + return { + ...actual, + constants: { ...actual.constants }, + lstatSync: vi.fn(actual.lstatSync), + openSync: vi.fn(actual.openSync), + fstatSync: vi.fn(actual.fstatSync), + readFileSync: vi.fn(actual.readFileSync), + closeSync: vi.fn(actual.closeSync), + }; +}); + +const actual = await vi.importActual('node:fs'); +const stagingFile = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46'; +const receipt = { + embeddingCheckpoint: { + kind: 'interrupted', + at: '2026-10-03T12:00:00.000Z', + nodesProcessed: 1, + totalNodes: 2, + chunksProcessed: 1, + model: 'test-model', + dimensions: 2, + provider: 'local', + recovery: { stagingFile, schemaFingerprint: 'test-schema', unsafeNodeIds: [] }, + }, +}; + +let dir: string; +let metadataPath: string; +beforeEach(() => { + vi.mocked(fs.lstatSync).mockImplementation(actual.lstatSync); + vi.mocked(fs.openSync).mockImplementation(actual.openSync); + vi.mocked(fs.fstatSync).mockImplementation(actual.fstatSync); + vi.mocked(fs.readFileSync).mockImplementation(actual.readFileSync); + vi.mocked(fs.closeSync).mockImplementation(actual.closeSync); + Object.assign(fs.constants, actual.constants); + vi.clearAllMocks(); + dir = actual.mkdtempSync(path.join(os.tmpdir(), 'gnx-recovery-metadata-race-')); + metadataPath = path.join(dir, 'gitnexus.json'); + actual.writeFileSync(path.join(dir, stagingFile), 'stage'); +}); +afterEach(() => actual.rmSync(dir, { recursive: true, force: true })); + +describe('readEmbeddingRecovery metadata races', () => { + it('does not read replacement metadata after checking the original file', () => { + actual.writeFileSync(metadataPath, JSON.stringify({ embeddingCheckpoint: null })); + let replaced = false; + vi.mocked(fs.lstatSync).mockImplementation((...args) => { + const stat = actual.lstatSync(...args); + if (args[0] === metadataPath && !replaced) { + replaced = true; + actual.renameSync(metadataPath, path.join(dir, 'original-metadata.json')); + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + } + return stat; + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(replaced).toBe(true); + const descriptor = vi.mocked(fs.openSync).mock.results[0]?.value; + expect(typeof descriptor).toBe('number'); + expect(fs.readFileSync).toHaveBeenCalledWith(descriptor, 'utf8'); + expect(fs.closeSync).toHaveBeenCalledWith(descriptor); + }); + + it('rejects a file replaced between opening and checking its identity', () => { + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + let replaced = false; + vi.mocked(fs.openSync).mockImplementation((...args) => { + const descriptor = actual.openSync(...args); + if (args[0] === metadataPath && !replaced) { + replaced = true; + actual.renameSync(metadataPath, path.join(dir, 'original-metadata.json')); + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + } + return descriptor; + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(replaced).toBe(true); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('reads regular metadata when no-follow opens are unavailable', () => { + Object.assign(fs.constants, { O_NOFOLLOW: 0, O_NONBLOCK: 0 }); + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + + expect(readEmbeddingRecovery(dir)?.stagingFile).toBe(stagingFile); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('refuses unverifiable file identity when no-follow opens are unavailable', () => { + Object.assign(fs.constants, { O_NOFOLLOW: 0, O_NONBLOCK: 0 }); + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + vi.mocked(fs.fstatSync).mockImplementation((...args) => { + const stat = actual.fstatSync(...args); + Object.defineProperty(stat, 'ino', { value: 0n }); + return stat; + }); + vi.mocked(fs.lstatSync).mockImplementation((...args) => { + const stat = actual.lstatSync(...args); + Object.defineProperty(stat, 'ino', { value: 0n }); + return stat; + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('rejects a symlink introduced before opening without no-follow support', () => { + Object.assign(fs.constants, { O_NOFOLLOW: 0, O_NONBLOCK: 0 }); + actual.writeFileSync(metadataPath, JSON.stringify({ embeddingCheckpoint: null })); + const target = path.join(dir, 'foreign-metadata.json'); + actual.writeFileSync(target, JSON.stringify(receipt)); + let replaced = false; + vi.mocked(fs.openSync).mockImplementation((...args) => { + if (args[0] === metadataPath && !replaced) { + replaced = true; + actual.rmSync(metadataPath); + actual.symlinkSync(target, metadataPath); + } + return actual.openSync(...args); + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(replaced).toBe(true); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('never falls back from a dangling primary symlink without no-follow support', () => { + Object.assign(fs.constants, { O_NOFOLLOW: 0, O_NONBLOCK: 0 }); + actual.symlinkSync(path.join(dir, 'missing-metadata.json'), metadataPath); + actual.writeFileSync(path.join(dir, 'meta.json'), JSON.stringify(receipt)); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(fs.readFileSync).not.toHaveBeenCalled(); + }); + + it('rejects a symlinked legacy receipt when the primary is absent', () => { + Object.assign(fs.constants, { O_NOFOLLOW: 0, O_NONBLOCK: 0 }); + const target = path.join(dir, 'foreign-metadata.json'); + actual.writeFileSync(target, JSON.stringify(receipt)); + actual.symlinkSync(target, path.join(dir, 'meta.json')); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('closes a descriptor when fstat fails without falling back to legacy metadata', () => { + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + actual.writeFileSync(path.join(dir, 'meta.json'), JSON.stringify(receipt)); + vi.mocked(fs.fstatSync).mockImplementationOnce(() => { + throw Object.assign(new Error('stat failed'), { code: 'EIO' }); + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it('closes a descriptor when JSON parsing fails without falling back', () => { + actual.writeFileSync(metadataPath, '{'); + actual.writeFileSync(path.join(dir, 'meta.json'), JSON.stringify(receipt)); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); + + it.skipIf(process.platform === 'win32')('rejects a substituted FIFO without blocking', () => { + actual.writeFileSync(metadataPath, JSON.stringify(receipt)); + let replaced = false; + vi.mocked(fs.openSync).mockImplementation((...args) => { + if (args[0] === metadataPath && !replaced) { + replaced = true; + actual.rmSync(metadataPath); + execFileSync('mkfifo', [metadataPath]); + } + return actual.openSync(...args); + }); + + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + expect(replaced).toBe(true); + expect(fs.readFileSync).not.toHaveBeenCalled(); + expect(fs.closeSync).toHaveBeenCalledOnce(); + }); +}); diff --git a/gitnexus/test/unit/embedding-recovery.test.ts b/gitnexus/test/unit/embedding-recovery.test.ts new file mode 100644 index 000000000..a8f1064a0 --- /dev/null +++ b/gitnexus/test/unit/embedding-recovery.test.ts @@ -0,0 +1,175 @@ +import { afterEach, beforeEach, describe, expect, it } from 'vitest'; +import { mkdtempSync, mkdirSync, rmSync, symlinkSync, writeFileSync } from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { + readEmbeddingRecovery, + resolveEmbeddingRecovery, +} from '../../src/storage/embedding-recovery.js'; + +const stagingFile = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46'; +const familySuffixes = [ + '', + '.wal', + '.shadow', + '.wal.checkpoint', + '.lock', + '.checkpoint.intent.lock', + '.checkpoint.apply.lock', +]; +const checkpoint = () => ({ + kind: 'interrupted', + at: '2026-10-03T12:00:00.000Z', + nodesProcessed: 1, + totalNodes: 3, + chunksProcessed: 2, + model: 'test-model', + dimensions: 2, + provider: 'local', + pendingNodeIds: ['active'], + recovery: { + stagingFile, + schemaFingerprint: 'test-schema', + unsafeNodeIds: ['active', 'incomplete-restore'], + }, +}); + +let dir: string; +beforeEach(() => { + dir = mkdtempSync(path.join(os.tmpdir(), 'gnx-embedding-recovery-')); + writeFileSync(path.join(dir, stagingFile), 'stage'); +}); +afterEach(() => rmSync(dir, { recursive: true, force: true })); + +describe('resolveEmbeddingRecovery', () => { + it('resolves the exact staged generation and carries all unsafe node IDs', () => { + expect(resolveEmbeddingRecovery(dir, checkpoint())).toEqual({ + ...checkpoint().recovery, + dbPath: path.join(dir, stagingFile), + familyFiles: familySuffixes.map((suffix) => stagingFile + suffix), + }); + }); + + it.each([ + '../lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46', + '/tmp/lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46', + 'branches/foreign/lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46', + '..\\lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46', + 'lbug', + 'lbug.new', + 'lbug.staging.orphan', + `${stagingFile}.wal`, + 'lbug.staging.00000000-0000-4000-8000-000000000000', + ])('rejects a foreign, unrecognized or missing generation: %s', (filename) => { + const marker = checkpoint(); + marker.recovery.stagingFile = filename; + expect(resolveEmbeddingRecovery(dir, marker)).toBeUndefined(); + }); + + it.each([ + null, + [], + { recovery: checkpoint().recovery }, + { ...checkpoint(), kind: 'partial' }, + { ...checkpoint(), kind: 'unverified-count' }, + { ...checkpoint(), kind: 'unknown' }, + { ...checkpoint(), at: 'invalid-time' }, + { ...checkpoint(), nodesProcessed: -1 }, + { ...checkpoint(), nodesProcessed: 4 }, + { ...checkpoint(), totalNodes: 1.5 }, + { ...checkpoint(), chunksProcessed: Number.NaN }, + { ...checkpoint(), model: '' }, + { ...checkpoint(), provider: null }, + { ...checkpoint(), dimensions: 0 }, + { ...checkpoint(), pendingNodeIds: [42] }, + { ...checkpoint(), pendingNodeIds: ['unexcluded'] }, + { ...checkpoint(), recovery: { ...checkpoint().recovery, schemaFingerprint: '' } }, + { ...checkpoint(), recovery: { ...checkpoint().recovery, unsafeNodeIds: undefined } }, + { ...checkpoint(), recovery: { ...checkpoint().recovery, unsafeNodeIds: [''] } }, + ])('rejects malformed or incompatible checkpoint shape %#', (marker) => { + expect(resolveEmbeddingRecovery(dir, marker)).toBeUndefined(); + }); + + it('accepts a durable restored stage even before a new embedding window completes', () => { + expect( + resolveEmbeddingRecovery(dir, { + ...checkpoint(), + nodesProcessed: 0, + chunksProcessed: 0, + pendingNodeIds: [], + }), + ).toBeDefined(); + }); + + it('rejects a directory in place of the staged database', () => { + rmSync(path.join(dir, stagingFile)); + mkdirSync(path.join(dir, stagingFile)); + expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined(); + }); + + it.each(familySuffixes)('rejects a symlink in the staged family: %s', (suffix) => { + const target = path.join(dir, 'external-file'); + writeFileSync(target, 'external'); + const candidate = path.join(dir, stagingFile + suffix); + rmSync(candidate, { force: true }); + symlinkSync(target, candidate); + expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined(); + }); + + it.each(familySuffixes.slice(1))('rejects a dangling staged sidecar symlink: %s', (suffix) => { + symlinkSync(path.join(dir, 'missing'), path.join(dir, stagingFile + suffix)); + expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined(); + }); + + it.each(familySuffixes.slice(1))('rejects a non-file staged sidecar: %s', (suffix) => { + mkdirSync(path.join(dir, stagingFile + suffix)); + expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined(); + }); +}); + +describe('readEmbeddingRecovery', () => { + const writeMeta = (filename: string, marker: unknown = checkpoint()): void => { + writeFileSync(path.join(dir, filename), JSON.stringify({ embeddingCheckpoint: marker })); + }; + + it('prefers the primary metadata file over a stale legacy reference', () => { + writeMeta('gitnexus.json'); + writeMeta('meta.json', null); + expect(readEmbeddingRecovery(dir)?.stagingFile).toBe(stagingFile); + }); + + it('loads the legacy mirror only when the primary metadata file is absent', () => { + writeMeta('meta.json'); + expect(readEmbeddingRecovery(dir)?.stagingFile).toBe(stagingFile); + }); + + it('does not resurrect a legacy reference when the primary file is malformed', () => { + writeFileSync(path.join(dir, 'gitnexus.json'), '{'); + writeMeta('meta.json'); + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + }); + + it('does not fall back from valid primary metadata with no recovery reference', () => { + writeMeta('gitnexus.json', null); + writeMeta('meta.json'); + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + }); + + it('rejects symlinked primary metadata rather than reading a foreign receipt', () => { + writeMeta('meta.json'); + symlinkSync(path.join(dir, 'meta.json'), path.join(dir, 'gitnexus.json')); + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + }); + + it('resolves branch-slot provenance within that slot despite a flat storagePath', () => { + const branchSlot = path.join(dir, 'branches', 'feature'); + mkdirSync(branchSlot, { recursive: true }); + writeFileSync(path.join(branchSlot, stagingFile), 'branch-stage'); + writeFileSync( + path.join(branchSlot, 'gitnexus.json'), + JSON.stringify({ storagePath: dir, embeddingCheckpoint: checkpoint() }), + ); + expect(readEmbeddingRecovery(branchSlot)?.dbPath).toBe(path.join(branchSlot, stagingFile)); + expect(readEmbeddingRecovery(dir)).toBeUndefined(); + }); +}); diff --git a/gitnexus/test/unit/embeddings-sync-command.test.ts b/gitnexus/test/unit/embeddings-sync-command.test.ts index 01afa40c3..ec33fee6d 100644 --- a/gitnexus/test/unit/embeddings-sync-command.test.ts +++ b/gitnexus/test/unit/embeddings-sync-command.test.ts @@ -3,13 +3,15 @@ * index lock, missing-DB preflight, identity fail-closed, tri-state count, * closeLbug masking, and hash-only cache load. */ -import { mkdtemp, mkdir, rm, writeFile } from 'node:fs/promises'; +import { existsSync } from 'node:fs'; +import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import path from 'node:path'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; const { acquireIndexLockMock, + ensurePrivateSharedGraphMock, releaseMock, getStoragePathsMock, loadMetaMock, @@ -27,6 +29,7 @@ const { reapEmbeddingSidecarMock, } = vi.hoisted(() => ({ acquireIndexLockMock: vi.fn(), + ensurePrivateSharedGraphMock: vi.fn(), releaseMock: vi.fn(), getStoragePathsMock: vi.fn(), loadMetaMock: vi.fn(), @@ -48,6 +51,10 @@ vi.mock('../../src/storage/git.js', () => ({ getGitRoot: () => '/tmp/emb-sync-repo', })); +vi.mock('../../src/core/shared-store-analyze.js', () => ({ + ensurePrivateSharedGraph: (...args: unknown[]) => ensurePrivateSharedGraphMock(...args), +})); + vi.mock('../../src/storage/index-lock.js', async (importOriginal) => ({ ...(await importOriginal()), acquireIndexLock: (...args: unknown[]) => acquireIndexLockMock(...args), @@ -113,6 +120,25 @@ async function run(inputPath = '/tmp/emb-sync-repo') { await embeddingsSyncCommand(inputPath); } +async function useRealIndexLock() { + vi.stubEnv('GITNEXUS_INDEX_LOCK_BACKEND', 'file'); + const actual = await vi.importActual( + '../../src/storage/index-lock.js', + ); + acquireIndexLockMock.mockImplementation( + async (...args: Parameters) => { + const lock = await actual.acquireIndexLock(...args); + return { + ...lock, + release: () => { + lock.release(); + releaseMock(); + }, + }; + }, + ); +} + describe('embeddingsSyncCommand writer safety (#3065)', () => { const tmpDirs: string[] = []; const originalEmbeddingUrl = process.env.GITNEXUS_EMBEDDING_URL; @@ -132,6 +158,7 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { beforeEach(() => { vi.resetModules(); acquireIndexLockMock.mockReset().mockResolvedValue(lockHandle()); + ensurePrivateSharedGraphMock.mockReset().mockResolvedValue(true); releaseMock.mockReset(); getStoragePathsMock.mockReset(); loadMetaMock.mockReset().mockResolvedValue({ ...BASE_META }); @@ -156,6 +183,7 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { }); afterEach(async () => { + vi.unstubAllEnvs(); if (originalEmbeddingUrl === undefined) delete process.env.GITNEXUS_EMBEDDING_URL; else process.env.GITNEXUS_EMBEDDING_URL = originalEmbeddingUrl; if (originalEmbeddingModel === undefined) delete process.env.GITNEXUS_EMBEDDING_MODEL; @@ -183,7 +211,7 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { await run(); - expect(acquireIndexLockMock).toHaveBeenCalledWith(dir); + expect(acquireIndexLockMock).toHaveBeenCalledWith(dir, { sweep: false }); expect(order[0]).toBe('lock'); expect(order.indexOf('loadMeta')).toBeGreaterThan(order.indexOf('lock')); expect(order.indexOf('init')).toBeGreaterThan(order.indexOf('loadMeta')); @@ -299,6 +327,90 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => { expect(releaseMock).toHaveBeenCalled(); }); + it.each([ + { + name: 'valid staged receipt', + recovery: { + stagingFile: 'lbug.staging.12345678-1234-4123-8123-123456789abc', + schemaFingerprint: 'test-schema', + unsafeNodeIds: ['n2', 'inherited-window-node'], + }, + }, + { name: 'null receipt', recovery: null }, + { name: 'malformed receipt', recovery: { stagingFile: 'invalid' } }, + { name: 'false receipt', recovery: false }, + ])('preserves a $name before any writable sync work', async ({ recovery }) => { + const { dir, lbugPath, metaPath } = await store(); + await useRealIndexLock(); + const lockPath = path.join(dir, 'analyze.lock'); + const sourcePath = path.join(dir, 'lbug.staging.12345678-1234-4123-8123-123456789abc'); + await writeFile(sourcePath, 'completed paid vectors'); + await writeFile(`${sourcePath}.wal`, 'unfinished window'); + const metadataBytes = JSON.stringify({ + ...BASE_META, + embeddingCheckpoint: { + ...IDENTITY, + at: '2026-01-01T00:00:00.000Z', + nodesProcessed: 1, + totalNodes: 2, + chunksProcessed: 1, + kind: 'interrupted', + pendingNodeIds: ['n2'], + recovery, + }, + }); + await writeFile(metaPath, metadataBytes); + loadMetaMock.mockImplementation(async () => { + expect(existsSync(lockPath)).toBe(true); + return JSON.parse(await readFile(metaPath, 'utf8')); + }); + resolveEmbeddingRuntimeMock.mockReturnValue(null); + + await expect(run()).rejects.toThrow( + /staged embeddings.*Run `gitnexus analyze` to recover them first/, + ); + + expect(acquireIndexLockMock).toHaveBeenCalledWith(dir, { sweep: false }); + expect(loadMetaMock.mock.invocationCallOrder[0]).toBeGreaterThan( + acquireIndexLockMock.mock.invocationCallOrder[0]!, + ); + expect(ensurePrivateSharedGraphMock).not.toHaveBeenCalled(); + expect(resolveEmbeddingIdentityMock).not.toHaveBeenCalled(); + expect(installEmbeddingRuntimeMock).not.toHaveBeenCalled(); + expect(initLbugMock).not.toHaveBeenCalled(); + expect(runEmbeddingPipelineMock).not.toHaveBeenCalled(); + expect(saveMetaMock).not.toHaveBeenCalled(); + expect(releaseMock).toHaveBeenCalledTimes(1); + expect(existsSync(lockPath)).toBe(false); + expect(await readFile(metaPath, 'utf8')).toBe(metadataBytes); + expect(await readFile(lbugPath, 'utf8')).toBe('db'); + expect(await readFile(sourcePath, 'utf8')).toBe('completed paid vectors'); + expect(await readFile(`${sourcePath}.wal`, 'utf8')).toBe('unfinished window'); + }); + + it('sweeps orphaned staging files before writable sync work without a recovery receipt', async () => { + const { dir, metaPath } = await store(); + await useRealIndexLock(); + await writeFile(metaPath, JSON.stringify(BASE_META)); + const sourcePath = path.join(dir, 'lbug.staging.12345678-1234-4123-8123-123456789abc'); + await writeFile(sourcePath, 'orphaned database'); + await writeFile(`${sourcePath}.wal`, 'orphaned WAL'); + loadMetaMock.mockImplementation(async () => JSON.parse(await readFile(metaPath, 'utf8'))); + ensurePrivateSharedGraphMock.mockImplementation(async () => { + expect(existsSync(path.join(dir, 'analyze.lock'))).toBe(true); + expect(existsSync(sourcePath)).toBe(false); + expect(existsSync(`${sourcePath}.wal`)).toBe(false); + return true; + }); + + await run(); + + expect(ensurePrivateSharedGraphMock).toHaveBeenCalledTimes(1); + expect(runEmbeddingPipelineMock).toHaveBeenCalledTimes(1); + expect(releaseMock).toHaveBeenCalledTimes(1); + expect(existsSync(path.join(dir, 'analyze.lock'))).toBe(false); + }); + it('persists an interrupted checkpoint from the pipeline checkpoint callbacks', async () => { // The resume contract lives in these callbacks; a mock that never invokes // them leaves the whole save path unexecuted. diff --git a/gitnexus/test/unit/index-lock.test.ts b/gitnexus/test/unit/index-lock.test.ts index 1a6846cc0..120e09727 100644 --- a/gitnexus/test/unit/index-lock.test.ts +++ b/gitnexus/test/unit/index-lock.test.ts @@ -17,6 +17,7 @@ import { existsSync, chmodSync, symlinkSync, + mkdirSync, } from 'node:fs'; import os from 'node:os'; import path from 'node:path'; @@ -178,6 +179,109 @@ describe('release', () => { }); describe('sweepStagingArtifacts', () => { + const recoveryStage = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46'; + const seedRecovery = (stagingFile = recoveryStage): void => { + writeFileSync( + path.join(dir, 'gitnexus.json'), + JSON.stringify({ + embeddingCheckpoint: { + kind: 'interrupted', + at: '2026-10-03T12:00:00.000Z', + nodesProcessed: 1, + totalNodes: 2, + chunksProcessed: 1, + model: 'test-model', + dimensions: 2, + provider: 'local', + pendingNodeIds: ['node-2'], + recovery: { + stagingFile, + schemaFingerprint: 'schema-v1', + unsafeNodeIds: ['node-2'], + }, + }, + }), + ); + }; + + it('retains the checkpoint-referenced family while reclaiming unrelated orphans', () => { + const family = [ + '', + '.wal', + '.shadow', + '.wal.checkpoint', + '.lock', + '.checkpoint.intent.lock', + '.checkpoint.apply.lock', + ].map((suffix) => recoveryStage + suffix); + for (const name of [...family, `${recoveryStage}.unexpected`, 'lbug.staging.orphan']) { + writeFileSync(path.join(dir, name), 'x'); + } + seedRecovery(); + + sweepStagingArtifacts(dir); + + for (const name of family) expect(existsSync(path.join(dir, name))).toBe(true); + expect(existsSync(path.join(dir, `${recoveryStage}.unexpected`))).toBe(false); + expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(false); + }); + + it('preserves a referenced generation when a shared caller acquires the lock', async () => { + writeFileSync(path.join(dir, recoveryStage), 'x'); + seedRecovery(); + const lock = await acquireIndexLock(dir); + try { + expect(existsSync(path.join(dir, recoveryStage))).toBe(true); + } finally { + lock.release(); + } + }); + + it('retains only the referenced family in a branch sub-slot', async () => { + const branchDir = path.join(dir, 'branches', 'feature'); + mkdirSync(branchDir, { recursive: true }); + seedRecovery(); + const metadata = JSON.parse(readFileSync(path.join(dir, 'gitnexus.json'), 'utf8')); + writeFileSync( + path.join(branchDir, 'gitnexus.json'), + JSON.stringify({ ...metadata, storagePath: dir }), + ); + writeFileSync(path.join(branchDir, recoveryStage), 'stage'); + writeFileSync(path.join(branchDir, 'lbug.staging.orphan'), 'orphan'); + writeFileSync(path.join(dir, 'lbug.staging.orphan'), 'other-slot'); + + const lock = await acquireIndexLock(branchDir); + try { + expect(existsSync(path.join(branchDir, recoveryStage))).toBe(true); + expect(existsSync(path.join(branchDir, 'lbug.staging.orphan'))).toBe(false); + expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(true); + } finally { + lock.release(); + } + }); + + it('rejects a symlinked reference and never follows it while reclaiming the stage', () => { + const external = path.join(dir, 'foreign-database'); + writeFileSync(external, 'external'); + symlinkSync(external, path.join(dir, recoveryStage)); + seedRecovery(); + + sweepStagingArtifacts(dir); + + expect(existsSync(path.join(dir, recoveryStage))).toBe(false); + expect(readFileSync(external, 'utf8')).toBe('external'); + }); + + it('does not sweep any generation when acquisition explicitly defers cleanup', async () => { + writeFileSync(path.join(dir, 'lbug.staging.orphan'), 'x'); + const lock = await acquireIndexLock(dir, { sweep: false }); + try { + expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(true); + } finally { + lock.release(); + } + }); + it('removes only staging files, never the live index or its sidecars', () => { const files = [ 'lbug', diff --git a/gitnexus/test/unit/run-analyze-fts-repair.test.ts b/gitnexus/test/unit/run-analyze-fts-repair.test.ts index b6de84c6b..b60cd60f1 100644 --- a/gitnexus/test/unit/run-analyze-fts-repair.test.ts +++ b/gitnexus/test/unit/run-analyze-fts-repair.test.ts @@ -1,13 +1,18 @@ import { execSync } from 'child_process'; import fs from 'fs/promises'; -import { afterEach, describe, expect, it, vi, type Mock } from 'vitest'; +import { basename } from 'node:path'; +import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'; import { getStoragePaths, loadMeta, saveMeta, type RepoMeta, } from '../../src/storage/repo-manager.js'; -import { EMBEDDING_DIMS, STALE_HASH_SENTINEL } from '../../src/core/lbug/schema.js'; +import { + EMBEDDING_DIMS, + STALE_HASH_SENTINEL, + SCHEMA_FINGERPRINT, +} from '../../src/core/lbug/schema.js'; import { getIndexIncompleteReasons } from '../../src/core/index-freshness.js'; import type { EmbeddingPipelineOptions, @@ -2382,12 +2387,20 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => { incrementalInProgress: { phase: 'full-rebuild' }, stats: { embeddings: 7 }, }); + expect(snapshots.postWindow?.embeddingCheckpoint?.recovery).toMatchObject({ + stagingFile: expect.stringMatching(/^lbug\.staging\.[a-f0-9-]+$/), + schemaFingerprint: SCHEMA_FINGERPRINT, + unsafeNodeIds: [], + }); // ── Window 2: the published count remains unchanged ──────────────── expect(snapshots.secondWindow).toMatchObject({ lastCommit: STALE_COMMIT, stats: { embeddings: 7 }, - embeddingCheckpoint: { pendingNodeIds: ['node-3', 'node-4'] }, + embeddingCheckpoint: { + pendingNodeIds: ['node-3', 'node-4'], + recovery: { unsafeNodeIds: ['node-3', 'node-4'] }, + }, }); // Only the finalize write — after the index is published — advances @@ -2427,6 +2440,14 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => { * NEXT run does with it). */ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => { + const actualPlatformDescriptor = Object.getOwnPropertyDescriptor(process, 'platform'); + + beforeEach(() => { + // These mocked checkpoint tests exercise the atomic rebuild used on POSIX. + // Select that same path on Windows; the in-place case below opts out. + vi.stubEnv('GITNEXUS_ATOMIC_WINDOWS_SWAP', '1'); + }); + const RESILIENCE_NODE_ID = 'Function:src/app.ts:handler:1'; const stubNode = { id: RESILIENCE_NODE_ID, @@ -2461,7 +2482,7 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => const mockResilienceHarness = ( controls: ResilienceControls, - ): { runEmbeddingPipeline: Mock; loadCachedEmbeddings: Mock } => { + ): { runEmbeddingPipeline: Mock; loadCachedEmbeddings: Mock; batchInsertEmbeddings: Mock } => { const loadCachedEmbeddings = vi.fn(async () => ({ embeddingNodeIds: new Set(), embeddings: [], @@ -2530,11 +2551,13 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => pipelineOptions: EmbeddingPipelineOptions, ): Promise => controls.pipeline(pipelineOptions), ); + const batchInsertEmbeddings = vi.fn(async () => undefined); vi.doMock('../../src/core/embeddings/embedding-pipeline.js', () => ({ runEmbeddingPipeline, + batchInsertEmbeddings, buildVectorIndex: vi.fn(async () => false), })); - return { runEmbeddingPipeline, loadCachedEmbeddings }; + return { runEmbeddingPipeline, loadCachedEmbeddings, batchInsertEmbeddings }; }; /** A checkpoint shaped exactly as `RepoMeta` declares it. */ @@ -2583,12 +2606,17 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => }; afterEach(() => { + if (actualPlatformDescriptor) { + Object.defineProperty(process, 'platform', actualPlatformDescriptor); + } vi.doUnmock('../../src/core/lbug/lbug-adapter.js'); vi.doUnmock('../../src/core/search/fts-indexes.js'); vi.doUnmock('../../src/core/ingestion/pipeline.js'); vi.doUnmock('../../src/storage/repo-manager.js'); vi.doUnmock('../../src/core/embeddings/embedding-identity.js'); vi.doUnmock('../../src/core/embeddings/embedding-pipeline.js'); + vi.doUnmock('../../src/core/embeddings/staged-embedding-recovery.js'); + vi.restoreAllMocks(); vi.resetModules(); vi.clearAllMocks(); vi.unstubAllEnvs(); @@ -3040,4 +3068,625 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => await tmpRepo.cleanup(); } }); + + const mockStagedFiles = async (): Promise => { + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => { + if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture'); + }); + vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => { + await fs.rm(dbPath, { force: true }); + }); + }; + + it.each([ + { mode: 'staged', atomicSwap: '1', checkpointCount: 7 }, + { mode: 'in-place', atomicSwap: '0', checkpointCount: 42 }, + ])( + 'keeps Windows $mode checkpoint counts consistent with the live index', + async ({ mode, atomicSwap, checkpointCount }) => { + const tmpRepo = await createTempDir('gitnexus-2790-checkpoint-meta-'); + try { + const { storagePath, lbugPath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { nodes: 2, embeddings: 7 } }); + await fs.writeFile(lbugPath, 'published fixture'); + const snapshots: Array = []; + mockResilienceHarness({ + count: [{ cnt: 42 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: [RESILIENCE_NODE_ID], + }); + snapshots.push(await loadMeta(storagePath)); + await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 3 }); + snapshots.push(await loadMeta(storagePath)); + return cleanResult(); + }, + }); + await mockStagedFiles(); + // Load modules on the actual host before changing only the platform + // branch exercised by runFullAnalysis; all native DB work is mocked. + const { runFullAnalysis } = await import('../../src/core/run-analyze.js'); + vi.stubEnv('GITNEXUS_ATOMIC_WINDOWS_SWAP', atomicSwap); + Object.defineProperty(process, 'platform', { value: 'win32', configurable: true }); + await runFullAnalysis( + tmpRepo.dbPath, + { force: true, embeddings: true, skipAgentsMd: true, skipSkills: true }, + { onProgress: () => {}, onLog: () => {} }, + ); + + expect(snapshots[0]?.stats?.embeddings).toBe(7); + expect(snapshots[1]?.stats?.embeddings).toBe(checkpointCount); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + const buildPath = vi.mocked(adapter.initLbug).mock.calls.at(-1)?.[0]; + if (mode === 'staged') { + expect(buildPath).toMatch(/\.staging\.[a-f0-9-]+$/); + expect(snapshots[1]?.embeddingCheckpoint?.recovery).toMatchObject({ + stagingFile: expect.stringMatching(/^lbug\.staging\.[a-f0-9-]+$/), + unsafeNodeIds: [], + }); + expect(await fs.readFile(lbugPath, 'utf8')).toBe('staged fixture'); + } else { + expect(buildPath).toBe(lbugPath); + expect(snapshots[0]?.embeddingCheckpoint?.recovery).toBeUndefined(); + expect(snapshots[1]?.embeddingCheckpoint?.recovery).toBeUndefined(); + } + const finalMeta = await loadMeta(storagePath); + expect(finalMeta?.stats?.embeddings).toBe(42); + expect(finalMeta?.embeddingCheckpoint).toBeUndefined(); + } finally { + if (actualPlatformDescriptor) { + Object.defineProperty(process, 'platform', actualPlatformDescriptor); + } + await tmpRepo.cleanup(); + } + }, + ); + + it.each(['close', 'rename'] as const)( + 'keeps the published count and database when %s fails before the staged publish', + async (failure) => { + const tmpRepo = await createTempDir('gitnexus-2790-checkpoint-meta-'); + try { + const { storagePath, lbugPath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { nodes: 2, embeddings: 7 } }); + await fs.writeFile(lbugPath, 'published fixture'); + const rename = fs.rename.bind(fs); + let checkpointMeta: RepoMeta | null = null; + let failedPublishRename: Mock | undefined; + mockResilienceHarness({ + count: [{ cnt: 42 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: [RESILIENCE_NODE_ID], + }); + await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 3 }); + checkpointMeta = await loadMeta(storagePath); + if (failure === 'close') { + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.closeLbug).mockRejectedValueOnce( + new Error('pre-publish close failed'), + ); + } else { + failedPublishRename = vi.fn(async () => { + throw Object.assign(new Error('staged publish rename failed'), { code: 'EIO' }); + }); + vi.spyOn(fs, 'rename').mockImplementation(async (source, destination) => { + if ( + String(source).startsWith(`${lbugPath}.staging.`) && + String(destination) === lbugPath + ) { + return failedPublishRename?.(); + } + return rename(source, destination); + }); + } + return cleanResult(); + }, + }); + await mockStagedFiles(); + expect( + await runAnalyze( + tmpRepo.dbPath, + { force: true, embeddings: true, skipAgentsMd: true, skipSkills: true }, + [], + ), + ).toMatchObject({ + message: + failure === 'close' ? 'pre-publish close failed' : 'staged publish rename failed', + }); + expect(checkpointMeta?.stats?.embeddings).toBe(7); + expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(7); + expect(await fs.readFile(lbugPath, 'utf8')).toBe('published fixture'); + const recovery = (await loadMeta(storagePath))?.embeddingCheckpoint?.recovery; + if (!recovery) throw new Error('expected durable stage after failed publish'); + expect(await fs.readFile(`${storagePath}/${recovery.stagingFile}`, 'utf8')).toBe( + 'staged fixture', + ); + if (failure === 'rename') expect(failedPublishRename).toHaveBeenCalledTimes(1); + } finally { + vi.restoreAllMocks(); + await tmpRepo.cleanup(); + } + }, + ); + + const seedRecovery = async (storagePath: string, repoPath: string) => { + const stagingFile = 'lbug.staging.11111111-1111-4111-8111-111111111111'; + const checkpoint = checkpointFixture({ + kind: 'interrupted', + pendingNodeIds: [], + recovery: { stagingFile, schemaFingerprint: SCHEMA_FINGERPRINT, unsafeNodeIds: [] }, + }); + await seedMeta(storagePath, repoPath, { + stats: { nodes: 2, embeddings: 7 }, + embeddingCheckpoint: checkpoint, + }); + await fs.writeFile(`${storagePath}/${stagingFile}`, 'previous durable source'); + vi.doMock('../../src/core/embeddings/staged-embedding-recovery.js', () => ({ + recoverStagedEmbeddings: vi.fn(async () => ({ rows: [], embeddingNodeIds: new Set() })), + })); + return { checkpoint, stagingFile }; + }; + + it.each( + ['1', '0'].flatMap((manualCheckpoint) => + ( + ['window-start crash', 'post-window crash', 'final-metadata failure', 'success'] as const + ).map((outcome) => ({ manualCheckpoint, outcome })), + ), + )( + 'preserves the in-place recovery receipt with manual checkpoints=$manualCheckpoint through $outcome', + async ({ manualCheckpoint, outcome }) => { + vi.stubEnv('GITNEXUS_WAL_MANUAL_CHECKPOINT', manualCheckpoint); + vi.stubEnv('GITNEXUS_INDEX_LOCK_BACKEND', 'file'); + const tmpRepo = await createTempDir('gitnexus-3456-opt-out-source-'); + try { + const { storagePath, lbugPath } = getStoragePaths(tmpRepo.dbPath); + const { checkpoint, stagingFile } = await seedRecovery(storagePath, tmpRepo.dbPath); + const original = { + ...checkpoint, + pendingNodeIds: ['original-pending'], + recovery: { + stagingFile, + schemaFingerprint: SCHEMA_FINGERPRINT, + unsafeNodeIds: ['original-pending', 'original-unsafe'], + }, + } satisfies NonNullable; + await seedMeta(storagePath, tmpRepo.dbPath, { + stats: { nodes: 2, embeddings: 7 }, + embeddingCheckpoint: original, + }); + const sourceFiles = [ + { filename: stagingFile, contents: 'previous durable source' }, + { filename: `${stagingFile}.wal`, contents: 'previous durable WAL' }, + { filename: `${stagingFile}.shadow`, contents: 'previous durable shadow' }, + ]; + for (const source of sourceFiles) { + await fs.writeFile(`${storagePath}/${source.filename}`, source.contents); + } + await fs.writeFile(lbugPath, 'previous published index'); + const { normalizeCachedEmbeddings } = + await import('../../src/core/embeddings/embedding-restore-spill.js'); + vi.doMock( + '../../src/core/embeddings/staged-embedding-recovery.js', + async (importActual) => ({ + ...(await importActual< + typeof import('../../src/core/embeddings/staged-embedding-recovery.js') + >()), + recoverStagedEmbeddings: vi.fn(async () => + normalizeCachedEmbeddings({ + embeddings: [ + { + nodeId: RESILIENCE_NODE_ID, + chunkIndex: 0, + startLine: 1, + endLine: 2, + contentHash: 'current-hash', + embedding: new Array(EMBEDDING_DIMS).fill(0), + }, + ], + }), + ), + }), + ); + const snapshots: Array = []; + const rename = fs.rename.bind(fs); + const { batchInsertEmbeddings } = mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['current-window'], + }); + snapshots.push(await loadMeta(storagePath)); + if (outcome === 'window-start crash') throw new Error(outcome); + await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 9 }); + snapshots.push(await loadMeta(storagePath)); + if (outcome === 'post-window crash') throw new Error(outcome); + if (outcome === 'final-metadata failure') { + vi.spyOn(fs, 'rename').mockImplementation(async (source, destination) => { + if (basename(String(destination)) === 'gitnexus.json') throw new Error(outcome); + return rename(source, destination); + }); + } + return cleanResult(); + }, + }); + await mockStagedFiles(); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.loadGraphToLbug).mockImplementation(async () => { + await fs.writeFile(lbugPath, 'in-place replacement'); + }); + // Import on the host first, then choose the Windows default in-place + // branch. The native adapter is mocked; source files and lock cleanup are real. + await import('../../src/core/run-analyze.js'); + vi.stubEnv('GITNEXUS_ATOMIC_WINDOWS_SWAP', '0'); + Object.defineProperty(process, 'platform', { value: 'win32', configurable: true }); + const error = await runAnalyze( + tmpRepo.dbPath, + { force: true, embeddings: true, skipAgentsMd: true, skipSkills: true }, + [], + ); + if (actualPlatformDescriptor) { + Object.defineProperty(process, 'platform', actualPlatformDescriptor); + } + expect(batchInsertEmbeddings).toHaveBeenCalled(); + expect(vi.mocked(adapter.initLbug).mock.calls.at(-1)?.[0]).toBe(lbugPath); + const finalMeta = await loadMeta(storagePath); + // Reacquiring the real lock performs the next retry's orphan sweep. + // A lost receipt would delete every byte of the proven old generation here. + const { acquireIndexLock } = await import('../../src/storage/index-lock.js'); + const lock = await acquireIndexLock(storagePath, { timeoutMs: 1000 }); + try { + if (outcome === 'success') { + expect(error).toBeNull(); + expect(finalMeta?.embeddingCheckpoint).toBeUndefined(); + expect(finalMeta?.stats?.embeddings).toBe(9); + expect(await fs.readFile(lbugPath, 'utf8')).toBe('in-place replacement'); + for (const source of sourceFiles) { + await expect(fs.stat(`${storagePath}/${source.filename}`)).rejects.toMatchObject({ + code: 'ENOENT', + }); + } + } else { + expect(error).toMatchObject({ message: outcome }); + for (const source of sourceFiles) { + expect(await fs.readFile(`${storagePath}/${source.filename}`, 'utf8')).toBe( + source.contents, + ); + } + expect(finalMeta?.embeddingCheckpoint).toEqual(original); + expect(finalMeta?.stats?.embeddings).toBe(outcome === 'window-start crash' ? 7 : 9); + } + expect(snapshots[0]?.embeddingCheckpoint).toEqual(original); + expect(snapshots[0]?.stats?.embeddings).toBe(7); + if (outcome !== 'window-start crash') { + expect(snapshots[1]?.embeddingCheckpoint).toEqual(original); + expect(snapshots[1]?.stats?.embeddings).toBe(9); + } + } finally { + lock.release(); + } + } finally { + if (actualPlatformDescriptor) { + Object.defineProperty(process, 'platform', actualPlatformDescriptor); + } + vi.restoreAllMocks(); + await tmpRepo.cleanup(); + } + }, + ); + + it('finishes a staged embedding run with manual checkpoints disabled', async () => { + vi.stubEnv('GITNEXUS_WAL_MANUAL_CHECKPOINT', '0'); + const tmpRepo = await createTempDir('gitnexus-3456-checkpoint-opt-out-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } }); + mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['active-node'], + }); + expect((await loadMeta(storagePath))?.embeddingCheckpoint?.recovery).toBeUndefined(); + await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 9 }); + const midRun = await loadMeta(storagePath); + expect(midRun?.embeddingCheckpoint?.recovery).toBeUndefined(); + expect(midRun?.stats?.embeddings).toBe(7); + return cleanResult(); + }, + }); + await mockStagedFiles(); + expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toBeNull(); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined(); + expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(9); + expect(await fs.readFile(getStoragePaths(tmpRepo.dbPath).lbugPath, 'utf8')).toBe( + 'staged fixture', + ); + } finally { + await tmpRepo.cleanup(); + } + }); + + it('keeps the previous durable source when manual checkpoints are disabled', async () => { + vi.stubEnv('GITNEXUS_WAL_MANUAL_CHECKPOINT', '0'); + const tmpRepo = await createTempDir('gitnexus-3456-opt-out-source-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + const { checkpoint: original, stagingFile } = await seedRecovery(storagePath, tmpRepo.dbPath); + mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['active-node'], + }); + await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 9 }); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original); + throw new Error('endpoint failure with manual checkpoints disabled'); + }, + }); + await mockStagedFiles(); + expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toMatchObject( + { message: 'endpoint failure with manual checkpoints disabled' }, + ); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original); + expect(await fs.readFile(`${storagePath}/${stagingFile}`, 'utf8')).toBe( + 'previous durable source', + ); + expect( + (await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')), + ).toEqual([stagingFile]); + } finally { + await tmpRepo.cleanup(); + } + }); + + it.skipIf(process.platform === 'win32')( + 'retains the referenced stage through a symlinked storage directory', + async () => { + const tmpRepo = await createTempDir('gitnexus-3456-storage-alias-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + const actualStorage = `${tmpRepo.dbPath}/actual-index`; + await fs.mkdir(actualStorage); + await fs.symlink(actualStorage, storagePath, 'dir'); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } }); + mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['active-node'], + }); + await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 3, chunksProcessed: 9 }); + throw new Error('endpoint failed after durable window'); + }, + }); + await mockStagedFiles(); + expect( + await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, []), + ).toMatchObject({ message: 'endpoint failed after durable window' }); + const recovery = (await loadMeta(storagePath))?.embeddingCheckpoint?.recovery; + if (!recovery) throw new Error('expected retained recovery generation'); + expect(await fs.readFile(`${actualStorage}/${recovery.stagingFile}`, 'utf8')).toBe( + 'staged fixture', + ); + expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(7); + } finally { + await tmpRepo.cleanup(); + } + }, + ); + + it.each(['checkpoint', 'metadata'] as const)( + 'keeps the previous source when %s fails before recovery handoff', + async (failure) => { + const tmpRepo = await createTempDir('gitnexus-3456-handoff-failure-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + const { checkpoint: original, stagingFile } = await seedRecovery( + storagePath, + tmpRepo.dbPath, + ); + const rename = fs.rename.bind(fs); + mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + if (failure === 'metadata') { + vi.spyOn(fs, 'rename').mockImplementation(async (source, destination) => { + if (basename(String(destination)) === 'gitnexus.json') + throw new Error('metadata write failed'); + return rename(source, destination); + }); + } else { + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.tryFlushWAL).mockResolvedValue(false); + } + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['active-node'], + }); + return cleanResult(); + }, + }); + await mockStagedFiles(); + const error = await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, []); + expect(error).toMatchObject({ + message: + failure === 'metadata' + ? 'metadata write failed' + : 'Could not checkpoint restored embeddings before recovery handoff.', + }); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original); + expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(7); + expect(await fs.readFile(`${storagePath}/${stagingFile}`, 'utf8')).toBe( + 'previous durable source', + ); + expect( + (await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')), + ).toEqual([stagingFile]); + } finally { + vi.restoreAllMocks(); + await tmpRepo.cleanup(); + } + }, + ); + + it('keeps an active window unsafe when its completion checkpoint fails', async () => { + const tmpRepo = await createTempDir('gitnexus-3456-completion-failure-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } }); + mockResilienceHarness({ + count: [{ cnt: 9 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 3, + chunksProcessed: 0, + nodeIds: ['active-node'], + }); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.tryFlushWAL).mockResolvedValue(false); + await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 3, chunksProcessed: 9 }); + return cleanResult(); + }, + }); + await mockStagedFiles(); + expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toMatchObject( + { message: 'Could not checkpoint the completed embedding window for recovery.' }, + ); + const meta = await loadMeta(storagePath); + expect(meta?.stats?.embeddings).toBe(7); + expect(meta?.embeddingCheckpoint?.pendingNodeIds).toEqual(['active-node']); + expect(meta?.embeddingCheckpoint?.recovery?.unsafeNodeIds).toEqual(['active-node']); + if (!meta?.embeddingCheckpoint?.recovery) throw new Error('expected retained active window'); + expect( + await fs.readFile( + `${storagePath}/${meta.embeddingCheckpoint.recovery.stagingFile}`, + 'utf8', + ), + ).toBe('staged fixture'); + } finally { + await tmpRepo.cleanup(); + } + }); + + it('retains a failed stage and keeps future incomplete restore groups unsafe', async () => { + const tmpRepo = await createTempDir('gitnexus-3456-unsafe-restore-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, { stats: { nodes: 2, embeddings: 1 } }); + let completedWindow: RepoMeta | null = null; + const { loadCachedEmbeddings, batchInsertEmbeddings } = mockResilienceHarness({ + count: [{ cnt: 5 }], + pipeline: async (options) => { + await options.onCheckpointWindowStart?.({ + nodesProcessed: 0, + totalNodes: 2, + chunksProcessed: 0, + nodeIds: ['earlier-window-node'], + }); + await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 2, chunksProcessed: 1 }); + completedWindow = await loadMeta(storagePath); + throw new Error('Maximum database size exceeded'); + }, + }); + loadCachedEmbeddings.mockResolvedValue({ + embeddingNodeIds: new Set([RESILIENCE_NODE_ID]), + embeddings: [ + { + nodeId: RESILIENCE_NODE_ID, + chunkIndex: 0, + startLine: 1, + endLine: 2, + contentHash: 'current-hash', + embedding: new Array(EMBEDDING_DIMS).fill(0), + }, + ], + }); + batchInsertEmbeddings.mockRejectedValue(new Error('restore batch partially inserted')); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => { + if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture'); + }); + vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => { + await fs.rm(dbPath, { force: true }); + }); + const logs: string[] = []; + expect( + await runAnalyze( + tmpRepo.dbPath, + { embeddings: true, force: true, skipAgentsMd: true, skipSkills: true }, + logs, + ), + ).toMatchObject({ message: 'Maximum database size exceeded' }); + expect(completedWindow?.embeddingCheckpoint?.recovery?.unsafeNodeIds).toEqual([ + RESILIENCE_NODE_ID, + ]); + expect(completedWindow?.stats?.embeddings).toBe(1); + const recovery = (await loadMeta(storagePath))?.embeddingCheckpoint?.recovery; + if (!recovery) throw new Error('expected retained recovery generation'); + expect(await fs.readFile(`${storagePath}/${recovery.stagingFile}`, 'utf8')).toBe( + 'staged fixture', + ); + expect(logs).toContainEqual(expect.stringContaining('GITNEXUS_LBUG_MAX_DB_SIZE')); + } finally { + await tmpRepo.cleanup(); + } + }); + + it('reclaims a failed current stage before any recovery checkpoint exists', async () => { + const tmpRepo = await createTempDir('gitnexus-3456-no-checkpoint-'); + try { + const { storagePath } = getStoragePaths(tmpRepo.dbPath); + await seedMeta(storagePath, tmpRepo.dbPath, {}); + mockResilienceHarness({ + count: [{ cnt: 0 }], + pipeline: async () => { + throw new Error('failed before first window'); + }, + }); + const adapter = await import('../../src/core/lbug/lbug-adapter.js'); + vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => { + if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture'); + }); + vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => { + await fs.rm(dbPath, { force: true }); + }); + expect( + await runAnalyze( + tmpRepo.dbPath, + { embeddings: true, force: true, skipAgentsMd: true, skipSkills: true }, + [], + ), + ).toMatchObject({ message: 'failed before first window' }); + expect( + (await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')), + ).toEqual([]); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined(); + } finally { + await tmpRepo.cleanup(); + } + }); }); diff --git a/gitnexus/test/unit/staged-embedding-recovery-child.test.ts b/gitnexus/test/unit/staged-embedding-recovery-child.test.ts new file mode 100644 index 000000000..58711493b --- /dev/null +++ b/gitnexus/test/unit/staged-embedding-recovery-child.test.ts @@ -0,0 +1,300 @@ +import fs from 'node:fs'; +import { EventEmitter } from 'node:events'; +import os from 'node:os'; +import path from 'node:path'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; + +const h = vi.hoisted(() => ({ + dbCtor: vi.fn(), + connCtor: vi.fn(), + dbClose: vi.fn<() => Promise>(), + connClose: vi.fn<() => Promise>(), + query: vi.fn(), + abortBuilder: vi.fn(), + spawn: vi.fn(), +})); + +vi.mock('node:child_process', () => ({ spawn: h.spawn })); + +vi.mock('@ladybugdb/core', () => { + class Database { + constructor(...args: unknown[]) { + h.dbCtor(...args); + } + close = h.dbClose; + } + class Connection { + constructor(db: unknown) { + h.connCtor(db); + } + query = h.query; + close = h.connClose; + } + return { default: { Database, Connection } }; +}); + +vi.mock('../../src/core/embeddings/embedding-restore-spill.js', async (importOriginal) => { + const actual = + await importOriginal(); + return { + ...actual, + abortCachedEmbeddingsBuilder: ( + ...args: Parameters + ) => { + h.abortBuilder(...args); + return actual.abortCachedEmbeddingsBuilder(...args); + }, + }; +}); + +describe('staged embedding recovery child native lifecycle', () => { + const suffixes = [ + '', + '.wal', + '.shadow', + '.wal.checkpoint', + '.lock', + '.checkpoint.intent.lock', + '.checkpoint.apply.lock', + ]; + let tmp: string; + let dbPath: string; + let exportDir: string; + let originalArgv: string[]; + let originalExitCode: typeof process.exitCode; + + beforeEach(() => { + vi.resetModules(); + vi.clearAllMocks(); + h.dbCtor.mockReset(); + h.connCtor.mockReset(); + h.dbClose.mockReset().mockResolvedValue(undefined); + h.connClose.mockReset().mockResolvedValue(undefined); + h.query.mockReset().mockResolvedValue({ + hasNext: vi.fn().mockResolvedValue(false), + getNext: vi.fn(), + close: vi.fn().mockResolvedValue(undefined), + }); + tmp = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-recovery-child-')); + dbPath = path.join(tmp, 'lbug.stage-test'); + exportDir = path.join(tmp, 'export'); + fs.writeFileSync(dbPath, 'mock native database'); + fs.writeFileSync(`${dbPath}.wal`, 'retained WAL'); + fs.mkdirSync(exportDir); + originalArgv = process.argv; + originalExitCode = process.exitCode; + process.argv = [process.execPath, 'staged-embedding-recovery-child', dbPath, exportDir, '2']; + process.exitCode = undefined; + vi.spyOn(process.stderr, 'write').mockImplementation(() => true); + }); + + afterEach(() => { + vi.useRealTimers(); + process.argv = originalArgv; + process.exitCode = originalExitCode; + vi.restoreAllMocks(); + fs.rmSync(tmp, { recursive: true, force: true }); + }); + + function sourceFamily() { + return Object.fromEntries( + suffixes.map((suffix) => [suffix, fs.readFileSync(dbPath + suffix, 'utf8')]), + ); + } + + function seedCompleteFamily() { + for (const suffix of suffixes) fs.writeFileSync(dbPath + suffix, `retained ${suffix}`); + return sourceFamily(); + } + + async function runRejectedChild(message: string): Promise { + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => { + expect(process.exitCode).toBe(1); + expect(process.stderr.write).toHaveBeenCalledWith(`${message}\n`); + }); + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + expect(fs.readFileSync(`${dbPath}.wal`, 'utf8')).toBe('retained WAL'); + } + + it('closes the opened database and aborts the builder when Connection construction fails', async () => { + h.connCtor.mockImplementation(() => { + throw new Error('connection constructor failed'); + }); + + await runRejectedChild('connection constructor failed'); + + expect(h.dbClose).toHaveBeenCalledOnce(); + expect(h.connClose).not.toHaveBeenCalled(); + expect(h.abortBuilder).toHaveBeenCalledOnce(); + expect(h.query).not.toHaveBeenCalled(); + expect(h.dbCtor.mock.calls[0][7]).toBe(true); + }); + + it('aborts the builder when Database construction fails', async () => { + h.dbCtor.mockImplementation(() => { + throw new Error('database constructor failed'); + }); + + await runRejectedChild('database constructor failed'); + + expect(h.abortBuilder).toHaveBeenCalledOnce(); + expect(h.dbClose).not.toHaveBeenCalled(); + expect(h.connCtor).not.toHaveBeenCalled(); + }); + + it('rejects output and closes the database when Connection close fails', async () => { + h.connClose.mockRejectedValue(new Error('connection close failed')); + + await runRejectedChild('connection close failed'); + + expect(h.dbClose).toHaveBeenCalled(); + expect(h.abortBuilder).toHaveBeenCalledOnce(); + }); + + it('rejects output when Database close fails', async () => { + h.dbClose.mockRejectedValue(new Error('database close failed')); + + await runRejectedChild('database close failed'); + + expect(h.connClose).toHaveBeenCalled(); + expect(h.abortBuilder).toHaveBeenCalledOnce(); + }); + + it('confines writable replay and failed checkpoint close to a separate copied family', async () => { + const sourceBefore = seedCompleteFamily(); + h.dbCtor.mockImplementation((openedPath: string) => { + for (const suffix of suffixes) { + expect(fs.readFileSync(openedPath + suffix, 'utf8')).toBe(sourceBefore[suffix]); + fs.writeFileSync(openedPath + suffix, `replayed ${suffix}`); + } + }); + h.dbClose.mockImplementation(async () => { + const openedPath = h.dbCtor.mock.calls[0][0] as string; + fs.writeFileSync(openedPath, 'partial checkpoint'); + throw new Error('checkpoint close failed'); + }); + + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => expect(process.exitCode).toBe(1)); + + expect(h.dbCtor.mock.calls[0][0]).not.toBe(dbPath); + expect(path.relative(exportDir, h.dbCtor.mock.calls[0][0] as string)).not.toMatch(/^\.\./); + expect(sourceFamily()).toEqual(sourceBefore); + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + expect(process.stderr.write).toHaveBeenCalledWith('checkpoint close failed\n'); + }); + + it('reclaims a timed-out writer copy without changing the retained source', async () => { + const sourceBefore = seedCompleteFamily(); + h.query.mockReturnValue(new Promise(() => {})); + h.dbCtor.mockImplementation((openedPath: string) => { + for (const suffix of suffixes) fs.writeFileSync(openedPath + suffix, 'writer opened'); + }); + let childImport: Promise | undefined; + const child = Object.assign(new EventEmitter(), { + stderr: new EventEmitter(), + kill: vi.fn(() => { + queueMicrotask(() => child.emit('close', null, 'SIGKILL')); + return true; + }), + }); + h.spawn.mockImplementation((_command: string, args: string[]) => { + process.argv = [process.execPath, 'staged-embedding-recovery-child', ...args.slice(-3)]; + childImport = import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + return child; + }); + const { recoverStagedEmbeddings } = + await import('../../src/core/embeddings/staged-embedding-recovery.js'); + vi.useFakeTimers(); + const recovering = expect( + recoverStagedEmbeddings(dbPath, { dimensions: 2, timeoutMs: 500 }), + ).rejects.toThrow(/timeout/); + await childImport; + expect(h.dbCtor).toHaveBeenCalledOnce(); + const openedPath = h.dbCtor.mock.calls[0][0] as string; + + await vi.advanceTimersByTimeAsync(500); + await recovering; + + expect(child.kill).toHaveBeenCalledWith('SIGKILL'); + expect(openedPath).not.toBe(dbPath); + expect(fs.existsSync(path.dirname(openedPath))).toBe(false); + expect(sourceFamily()).toEqual(sourceBefore); + expect(h.dbClose).not.toHaveBeenCalled(); + }); + + it.each(['.wal', '.checkpoint.intent.lock', '.checkpoint.apply.lock'])( + 'refuses a dangling family symlink before native open: %s', + async (suffix) => { + fs.rmSync(dbPath + suffix, { force: true }); + fs.symlinkSync(path.join(tmp, 'missing-sidecar'), dbPath + suffix); + + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => expect(process.exitCode).toBe(1)); + + expect(h.dbCtor).not.toHaveBeenCalled(); + expect(process.stderr.write).toHaveBeenCalledWith( + 'staged embedding family is not a regular file\n', + ); + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + }, + ); + + it('refuses a family entry replaced with a symlink between lstat and open', async () => { + const foreignPath = path.join(tmp, 'foreign-file'); + fs.writeFileSync(foreignPath, 'foreign'); + const open = fs.openSync; + const read = vi.spyOn(fs, 'readSync'); + vi.spyOn(fs, 'openSync').mockImplementation((...args) => { + if (args[0] === dbPath) { + fs.rmSync(dbPath); + fs.symlinkSync(foreignPath, dbPath); + } + return open(...args); + }); + + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => expect(process.exitCode).toBe(1)); + + expect(read).not.toHaveBeenCalled(); + expect(h.dbCtor).not.toHaveBeenCalled(); + expect(fs.readFileSync(foreignPath, 'utf8')).toBe('foreign'); + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + }); + + it('fails closed when copying the complete family runs out of space', async () => { + const sourceBefore = seedCompleteFamily(); + vi.spyOn(fs, 'writeFileSync').mockImplementation(() => { + throw Object.assign(new Error('copy ran out of space'), { code: 'ENOSPC' }); + }); + + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => expect(process.exitCode).toBe(1)); + + expect(h.dbCtor).not.toHaveBeenCalled(); + expect(sourceFamily()).toEqual(sourceBefore); + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + expect(process.stderr.write).toHaveBeenCalledWith('copy ran out of space\n'); + }); + + it('writes the manifest only after both native closes succeed', async () => { + const closed: string[] = []; + h.connClose.mockImplementation(async () => { + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + closed.push('connection'); + }); + h.dbClose.mockImplementation(async () => { + expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false); + closed.push('database'); + }); + + await import('../../src/core/embeddings/staged-embedding-recovery-child.js'); + await vi.waitFor(() => expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(true)); + + expect(closed).toEqual(['connection', 'database']); + expect(h.abortBuilder).not.toHaveBeenCalled(); + expect(process.exitCode).toBeUndefined(); + expect(fs.readFileSync(`${dbPath}.wal`, 'utf8')).toBe('retained WAL'); + }); +}); diff --git a/gitnexus/test/unit/staged-embedding-recovery.test.ts b/gitnexus/test/unit/staged-embedding-recovery.test.ts new file mode 100644 index 000000000..552e48dc5 --- /dev/null +++ b/gitnexus/test/unit/staged-embedding-recovery.test.ts @@ -0,0 +1,95 @@ +import fs from 'node:fs'; +import os from 'node:os'; +import path from 'node:path'; +import { afterEach, describe, expect, it } from 'vitest'; +import { + createCachedEmbeddingsBuilder, + disposeEmbeddingSpill, + finalizeCachedEmbeddingsSnapshot, + ingestCachedEmbeddingRow, + materializeCachedEmbeddings, + type CachedEmbeddingsSnapshot, +} from '../../src/core/embeddings/embedding-restore-spill.js'; +import { + mergeRecoveredEmbeddings, + validateRecoveredNodeGroups, +} from '../../src/core/embeddings/staged-embedding-recovery.js'; + +describe('staged embedding recovery', () => { + const snapshots: CachedEmbeddingsSnapshot[] = []; + let tmp: string | undefined; + afterEach(() => { + for (const snapshot of snapshots) disposeEmbeddingSpill(snapshot.spill); + snapshots.length = 0; + if (tmp) fs.rmSync(tmp, { recursive: true, force: true }); + tmp = undefined; + }); + + function snapshot( + rows: { nodeId: string; chunkIndex: number; contentHash?: string; embedding?: number[] }[], + ) { + tmp ??= fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-test-')); + const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0, spillDir: tmp }); + for (const row of rows) { + ingestCachedEmbeddingRow( + builder, + { startLine: 1, endLine: 2, embedding: [1, 2], ...row }, + true, + ); + } + const result = finalizeCachedEmbeddingsSnapshot(builder); + snapshots.push(result); + return result; + } + + it('accepts only complete groups with one content hash and unique contiguous chunk ordinals', () => { + const cached = snapshot([ + { nodeId: 'complete', chunkIndex: 1, contentHash: 'same' }, + { nodeId: 'complete', chunkIndex: 0, contentHash: 'same' }, + { nodeId: 'gap', chunkIndex: 1, contentHash: 'same' }, + { nodeId: 'duplicate', chunkIndex: 0, contentHash: 'same' }, + { nodeId: 'duplicate', chunkIndex: 0, contentHash: 'same' }, + { nodeId: 'mixed', chunkIndex: 0, contentHash: 'old' }, + { nodeId: 'mixed', chunkIndex: 1, contentHash: 'new' }, + { nodeId: 'no-hash', chunkIndex: 0 }, + { nodeId: 'unsafe', chunkIndex: 0, contentHash: 'same' }, + ]); + expect([...validateRecoveredNodeGroups(cached.rows, new Set(['unsafe']))]).toEqual([ + 'complete', + ]); + }); + + it('replaces an entire published node group and keeps unrelated rows without retaining vector arrays', () => { + const live = snapshot([ + { nodeId: 'changed', chunkIndex: 0, contentHash: 'old' }, + { nodeId: 'changed', chunkIndex: 1, contentHash: 'old' }, + { nodeId: 'other', chunkIndex: 0, contentHash: 'other' }, + ]); + const recovered = snapshot([ + { nodeId: 'changed', chunkIndex: 0, contentHash: 'new', embedding: [7, 8] }, + ]); + const merged = mergeRecoveredEmbeddings(live, recovered); + snapshots.push(merged); + expect(merged.embeddings).toEqual([]); + expect(merged.rows).toHaveLength(2); + expect(materializeCachedEmbeddings(merged, merged.rows)).toEqual([ + expect.objectContaining({ nodeId: 'other', contentHash: 'other' }), + expect.objectContaining({ + nodeId: 'changed', + chunkIndex: 0, + contentHash: 'new', + embedding: [7, 8], + }), + ]); + expect(fs.existsSync(live.spill.path)).toBe(true); + expect(fs.existsSync(recovered.spill.path)).toBe(true); + }); + + it('does not accept missing vector bytes during a merge', () => { + const cached = snapshot([{ nodeId: 'complete', chunkIndex: 0, contentHash: 'same' }]); + fs.truncateSync(cached.spill.path, 12); + expect(() => mergeRecoveredEmbeddings(snapshot([]), cached)).toThrow( + /short embedding spill read/, + ); + }); +});