From c5e013a8d8fc8d061e4d1b387c2bc16675ab86e6 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Gerg=C5=91=20Magyar?= Date: Sat, 3 Oct 2026 15:58:09 +0100 Subject: [PATCH] test(embeddings): cover recovery failure boundaries --- gitnexus/test/unit/embedding-recovery.test.ts | 166 ++++++++ gitnexus/test/unit/index-lock.test.ts | 98 +++++ .../test/unit/run-analyze-fts-repair.test.ts | 358 +++++++++++++++++- .../unit/staged-embedding-recovery.test.ts | 95 +++++ 4 files changed, 710 insertions(+), 7 deletions(-) create mode 100644 gitnexus/test/unit/embedding-recovery.test.ts create mode 100644 gitnexus/test/unit/staged-embedding-recovery.test.ts diff --git a/gitnexus/test/unit/embedding-recovery.test.ts b/gitnexus/test/unit/embedding-recovery.test.ts new file mode 100644 index 000000000..a826f8238 --- /dev/null +++ b/gitnexus/test/unit/embedding-recovery.test.ts @@ -0,0 +1,166 @@ +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 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: ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'].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(['', '.wal', '.shadow', '.wal.checkpoint', '.lock'])( + '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('rejects a dangling staged sidecar symlink', () => { + symlinkSync(path.join(dir, 'missing'), path.join(dir, `${stagingFile}.wal`)); + 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/index-lock.test.ts b/gitnexus/test/unit/index-lock.test.ts index 1a6846cc0..40ed16891 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,103 @@ 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'].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 2f0304ef4..e150b9b23 100644 --- a/gitnexus/test/unit/run-analyze-fts-repair.test.ts +++ b/gitnexus/test/unit/run-analyze-fts-repair.test.ts @@ -7,7 +7,11 @@ import { 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, @@ -2375,19 +2379,27 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => { }); expect(snapshots.windowStart?.lastCommit).not.toBe(currentCommit); - // ── Post-window: the one save that legitimately measured the count ── + // The measured count belongs to the unpublished staging generation. expect(snapshots.postWindow).toMatchObject({ lastCommit: STALE_COMMIT, fileHashes: STALE_HASHES, incrementalInProgress: { phase: 'full-rebuild' }, - stats: { embeddings: LIVE_EMBEDDING_COUNT }, + stats: { embeddings: 7 }, + }); + expect(snapshots.postWindow?.embeddingCheckpoint?.recovery).toMatchObject({ + stagingFile: expect.stringMatching(/^lbug\.staging\.[a-f0-9-]+$/), + schemaFingerprint: SCHEMA_FINGERPRINT, + unsafeNodeIds: [], }); // ── Window 2: no stale restatement over the measured figure ──────── expect(snapshots.secondWindow).toMatchObject({ lastCommit: STALE_COMMIT, - stats: { embeddings: LIVE_EMBEDDING_COUNT }, - embeddingCheckpoint: { pendingNodeIds: ['node-3', 'node-4'] }, + stats: { embeddings: 7 }, + embeddingCheckpoint: { + pendingNodeIds: ['node-3', 'node-4'], + recovery: { unsafeNodeIds: ['node-3', 'node-4'] }, + }, }); // Only the finalize write — after the index is published — advances @@ -2461,7 +2473,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 +2542,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. */ @@ -2589,6 +2603,8 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => 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(); @@ -3039,4 +3055,332 @@ 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 }); + }); + }; + + 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('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(); + }, + }); + expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toBeNull(); + expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined(); + expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(9); + } 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 (String(destination).endsWith('/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.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/, + ); + }); +});