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/run-analyze-fts-repair.test.ts b/gitnexus/test/unit/run-analyze-fts-repair.test.ts index e150b9b23..1ec3a4d8f 100644 --- a/gitnexus/test/unit/run-analyze-fts-repair.test.ts +++ b/gitnexus/test/unit/run-analyze-fts-repair.test.ts @@ -1,6 +1,6 @@ import { execSync } from 'child_process'; import fs from 'fs/promises'; -import { afterEach, describe, expect, it, vi, type Mock } from 'vitest'; +import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest'; import { getStoragePaths, loadMeta, @@ -2219,10 +2219,11 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => { vi.unstubAllEnvs(); }); - it('preserves lastCommit / fileHashes / the dirty flag, and never restates a stale count', async () => { + it('keeps staging counts out of published metadata until the index is swapped', async () => { const STALE_COMMIT = '1111111111111111111111111111111111111111'; const STALE_HASHES = { 'src/app.ts': 'stale-hash' }; const LIVE_EMBEDDING_COUNT = 42; + vi.stubEnv('GITNEXUS_ATOMIC_WINDOWS_SWAP', '1'); vi.doMock('../../src/core/lbug/lbug-adapter.js', () => ({ initLbug: vi.fn(async () => undefined), @@ -2439,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, @@ -2597,6 +2606,9 @@ 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'); @@ -2715,6 +2727,7 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => */ it('carries the mid-run count forward, not the run-start snapshot, so --force still loads the cache', async () => { const MID_RUN_COUNT = 12; + vi.stubEnv('GITNEXUS_ATOMIC_WINDOWS_SWAP', '1'); const tmpRepo = await createTempDir('gitnexus-2790r-latest-meta-'); try { const { storagePath } = getStoragePaths(tmpRepo.dbPath); @@ -3066,6 +3079,144 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () => }); }; + 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({ 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..8e9255109 --- /dev/null +++ b/gitnexus/test/unit/staged-embedding-recovery-child.test.ts @@ -0,0 +1,158 @@ +import fs from 'node:fs'; +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(), +})); + +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', () => { + 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(() => { + process.argv = originalArgv; + process.exitCode = originalExitCode; + vi.restoreAllMocks(); + fs.rmSync(tmp, { recursive: true, force: true }); + }); + + 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('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'); + }); +});