diff --git a/gitnexus/test/unit/api-fts-mode.test.ts b/gitnexus/test/unit/api-fts-mode.test.ts index 632bd37d3..56e940b74 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,79 @@ beforeEach(() => { })); }); +describe('POST /api/embed staged recovery preflight', () => { + 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 }) => { + const metaPath = path.join(entry.storagePath, 'gitnexus.json'); + const sourcePath = path.join( + entry.storagePath, + 'lbug.staging.12345678-1234-4123-8123-123456789abc', + ); + 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 () => 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); + 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.readFileSync(metaPath, 'utf8')).toBe(metadataBytes); + 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)); + }); +}); + 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 +359,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/embeddings-sync-command.test.ts b/gitnexus/test/unit/embeddings-sync-command.test.ts index 01afa40c3..ac0a4f27b 100644 --- a/gitnexus/test/unit/embeddings-sync-command.test.ts +++ b/gitnexus/test/unit/embeddings-sync-command.test.ts @@ -3,13 +3,14 @@ * 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 { 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 +28,7 @@ const { reapEmbeddingSidecarMock, } = vi.hoisted(() => ({ acquireIndexLockMock: vi.fn(), + ensurePrivateSharedGraphMock: vi.fn(), releaseMock: vi.fn(), getStoragePathsMock: vi.fn(), loadMetaMock: vi.fn(), @@ -48,6 +50,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), @@ -132,6 +138,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 }); @@ -299,6 +306,60 @@ 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, metaPath } = await store(); + 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 () => 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); + 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(await readFile(metaPath, 'utf8')).toBe(metadataBytes); + expect(await readFile(sourcePath, 'utf8')).toBe('completed paid vectors'); + expect(await readFile(`${sourcePath}.wal`, 'utf8')).toBe('unfinished window'); + }); + 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/run-analyze-fts-repair.test.ts b/gitnexus/test/unit/run-analyze-fts-repair.test.ts index 5d616cf30..b60cd60f1 100644 --- a/gitnexus/test/unit/run-analyze-fts-repair.test.ts +++ b/gitnexus/test/unit/run-analyze-fts-repair.test.ts @@ -3235,6 +3235,157 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => 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-');