test(embeddings): retain staged receipts across writer retries

Address PR #3463 with eight in-place retry cases covering manual checkpoints on/off, interrupted windows, final metadata failure, and successful publication. Verify the original receipt and staged/WAL/shadow bytes survive actual lock reacquisition until publication.

Add eight server/embed and embeddings-sync cases covering recovery-bearing and malformed fields, no writable or provider work, unchanged fixture bytes, and lock release.

Validation: 222 unit tests and 15 native integration tests passed. Build, TypeScript, changed-file ESLint (warnings only), Prettier, and diff checks passed. Refreshed graph change analysis mapped all six changed files.
This commit is contained in:
Gergő Magyar 2026-10-03 20:41:15 +01:00 • committed by GitHub
parent 9d997b74c1
commit 1915460ae7
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
3 changed files with 318 additions and 2 deletions

View file

@ -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<typeof import('../../src/storage/repo-manager.js')>()),
loadMeta: mocks.loadMeta,
saveMeta: mocks.saveMeta,
listRegisteredRepos: mocks.listRegisteredRepos,
}));
vi.mock('../../src/storage/index-lock.js', async (importOriginal) => ({
...(await importOriginal<typeof import('../../src/storage/index-lock.js')>()),
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<typeof import('../../src/storage/storage-resolver.js')>()),
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<string, unknown> = {}) {
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);
},
);
});

View file

@ -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<typeof import('../../src/storage/index-lock.js')>()),
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.

View file

@ -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<RepoMeta['embeddingCheckpoint']>;
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<RepoMeta | null> = [];
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-');