test(embeddings): cover Windows counts and recovery failures

This commit is contained in:
Gergő Magyar 2026-10-03 18:03:27 +01:00 • committed by GitHub
parent c93cdbdebb
commit 390e18a2fe
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 510 additions and 2 deletions

View file

@ -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<typeof import('node:fs')>();
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<typeof import('node:fs')>('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();
});
});

View file

@ -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<RepoMeta | null> = [];
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({

View file

@ -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<void>>(),
connClose: vi.fn<() => Promise<void>>(),
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<typeof import('../../src/core/embeddings/embedding-restore-spill.js')>();
return {
...actual,
abortCachedEmbeddingsBuilder: (
...args: Parameters<typeof actual.abortCachedEmbeddingsBuilder>
) => {
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<void> {
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');
});
});