test(recovery): preserve sources across rejected writers (#3463)

Add real-lock coverage for CLI and HTTP recovery preflights, source-copy rejection and timeout coverage, and complete checkpoint-family assertions.

Validation for this review-fix series: build, scoped ESLint/Prettier, and 142 focused tests pass. The full suite completed with 23,743 passing tests, 42 failures, 100 skips, and one unexpected native worker exit. All 41 non-catalog failing cases reproduce on unchanged c3437c2. The untouched VECTOR-catalog case passes all 10 tests on both baseline and changed-code reruns. The interrupted incremental-orchestration suite is being rerun in isolation.

Note: pre-existing failures in test/unit/{auto-sync-runner,auto-sync,cli-index-help,cli-update-notice,file-lock,hook-db-lock-probe,index-lock,process-identity,update-check,worker-pool-resilience}.test.ts and test/integration/{cli,mcp}/update-notice.test.ts are not addressed by this PR.
This commit is contained in:
Gergő Magyar 2026-10-04 08:40:41 +01:00 • committed by GitHub
parent 6d6807e775
commit 1a3f0f96b4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
5 changed files with 299 additions and 25 deletions

View file

@ -154,6 +154,31 @@ beforeEach(() => {
});
describe('POST /api/embed staged recovery preflight', () => {
afterEach(() => {
vi.unstubAllEnvs();
fs.rmSync(entry.storagePath, { recursive: true, force: true });
fs.mkdirSync(entry.storagePath, { recursive: true });
});
async function useRealIndexLock() {
vi.stubEnv('GITNEXUS_INDEX_LOCK_BACKEND', 'file');
const actual = await vi.importActual<typeof import('../../src/storage/index-lock.js')>(
'../../src/storage/index-lock.js',
);
mocks.acquireIndexLock.mockImplementation(
async (...args: Parameters<typeof actual.acquireIndexLock>) => {
const lock = await actual.acquireIndexLock(...args);
return {
...lock,
release: () => {
lock.release();
mocks.releaseIndexLock();
},
};
},
);
}
it.each([
{
name: 'valid staged receipt',
@ -167,11 +192,15 @@ describe('POST /api/embed staged recovery preflight', () => {
{ name: 'malformed receipt', recovery: { stagingFile: 'invalid' } },
{ name: 'false receipt', recovery: false },
])('preserves a $name and releases the job locks', async ({ recovery }) => {
await useRealIndexLock();
const lockPath = path.join(entry.storagePath, 'analyze.lock');
const lbugPath = path.join(entry.storagePath, 'lbug');
const metaPath = path.join(entry.storagePath, 'gitnexus.json');
const sourcePath = path.join(
entry.storagePath,
'lbug.staging.12345678-1234-4123-8123-123456789abc',
);
fs.writeFileSync(lbugPath, 'published graph');
fs.writeFileSync(sourcePath, 'completed paid vectors');
fs.writeFileSync(`${sourcePath}.wal`, 'unfinished window');
const metadataBytes = JSON.stringify({
@ -193,7 +222,10 @@ describe('POST /api/embed staged recovery preflight', () => {
},
});
fs.writeFileSync(metaPath, metadataBytes);
mocks.loadMeta.mockImplementation(async () => JSON.parse(fs.readFileSync(metaPath, 'utf8')));
mocks.loadMeta.mockImplementation(async () => {
expect(fs.existsSync(lockPath)).toBe(true);
return JSON.parse(fs.readFileSync(metaPath, 'utf8'));
});
mocks.withLbugDb.mockResolvedValue(undefined);
await invoke('/api/embed');
@ -208,7 +240,7 @@ describe('POST /api/embed staged recovery preflight', () => {
),
}),
);
expect(mocks.acquireIndexLock).toHaveBeenCalledWith(entry.storagePath);
expect(mocks.acquireIndexLock).toHaveBeenCalledWith(entry.storagePath, { sweep: false });
expect(mocks.loadMeta.mock.invocationCallOrder[0]).toBeGreaterThan(
mocks.acquireIndexLock.mock.invocationCallOrder[0]!,
);
@ -216,13 +248,47 @@ describe('POST /api/embed staged recovery preflight', () => {
expect(mocks.withLbugDb).not.toHaveBeenCalled();
expect(mocks.runEmbeddingPipeline).not.toHaveBeenCalled();
expect(mocks.saveMeta).not.toHaveBeenCalled();
expect(fs.existsSync(lockPath)).toBe(false);
expect(fs.readFileSync(metaPath, 'utf8')).toBe(metadataBytes);
expect(fs.readFileSync(lbugPath, 'utf8')).toBe('published graph');
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));
expect(fs.existsSync(lockPath)).toBe(false);
});
it('sweeps orphaned staging files before a writable job without a recovery receipt', async () => {
await useRealIndexLock();
const metaPath = path.join(entry.storagePath, 'gitnexus.json');
fs.writeFileSync(metaPath, JSON.stringify({ repoPath: entry.path }));
const sourcePath = path.join(
entry.storagePath,
'lbug.staging.12345678-1234-4123-8123-123456789abc',
);
fs.writeFileSync(sourcePath, 'orphaned database');
fs.writeFileSync(`${sourcePath}.wal`, 'orphaned WAL');
mocks.loadMeta.mockImplementation(async () => JSON.parse(fs.readFileSync(metaPath, 'utf8')));
mocks.ensurePrivateSharedGraph.mockImplementation(async () => {
expect(fs.existsSync(path.join(entry.storagePath, 'analyze.lock'))).toBe(true);
expect(fs.existsSync(sourcePath)).toBe(false);
expect(fs.existsSync(`${sourcePath}.wal`)).toBe(false);
return true;
});
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: 'complete' }),
);
expect(mocks.ensurePrivateSharedGraph).toHaveBeenCalledTimes(1);
expect(mocks.withLbugDb).toHaveBeenCalledTimes(1);
expect(fs.existsSync(path.join(entry.storagePath, 'analyze.lock'))).toBe(false);
});
});

View file

@ -8,6 +8,15 @@ import {
} from '../../src/storage/embedding-recovery.js';
const stagingFile = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46';
const familySuffixes = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
];
const checkpoint = () => ({
kind: 'interrupted',
at: '2026-10-03T12:00:00.000Z',
@ -37,9 +46,7 @@ describe('resolveEmbeddingRecovery', () => {
expect(resolveEmbeddingRecovery(dir, checkpoint())).toEqual({
...checkpoint().recovery,
dbPath: path.join(dir, stagingFile),
familyFiles: ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'].map(
(suffix) => stagingFile + suffix,
),
familyFiles: familySuffixes.map((suffix) => stagingFile + suffix),
});
});
@ -100,20 +107,22 @@ describe('resolveEmbeddingRecovery', () => {
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.each(familySuffixes)('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`));
it.each(familySuffixes.slice(1))('rejects a dangling staged sidecar symlink: %s', (suffix) => {
symlinkSync(path.join(dir, 'missing'), path.join(dir, stagingFile + suffix));
expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined();
});
it.each(familySuffixes.slice(1))('rejects a non-file staged sidecar: %s', (suffix) => {
mkdirSync(path.join(dir, stagingFile + suffix));
expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined();
});
});

View file

@ -3,6 +3,7 @@
* index lock, missing-DB preflight, identity fail-closed, tri-state count,
* closeLbug masking, and hash-only cache load.
*/
import { existsSync } from 'node:fs';
import { mkdtemp, mkdir, readFile, rm, writeFile } from 'node:fs/promises';
import { tmpdir } from 'node:os';
import path from 'node:path';
@ -119,6 +120,25 @@ async function run(inputPath = '/tmp/emb-sync-repo') {
await embeddingsSyncCommand(inputPath);
}
async function useRealIndexLock() {
vi.stubEnv('GITNEXUS_INDEX_LOCK_BACKEND', 'file');
const actual = await vi.importActual<typeof import('../../src/storage/index-lock.js')>(
'../../src/storage/index-lock.js',
);
acquireIndexLockMock.mockImplementation(
async (...args: Parameters<typeof actual.acquireIndexLock>) => {
const lock = await actual.acquireIndexLock(...args);
return {
...lock,
release: () => {
lock.release();
releaseMock();
},
};
},
);
}
describe('embeddingsSyncCommand writer safety (#3065)', () => {
const tmpDirs: string[] = [];
const originalEmbeddingUrl = process.env.GITNEXUS_EMBEDDING_URL;
@ -163,6 +183,7 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => {
});
afterEach(async () => {
vi.unstubAllEnvs();
if (originalEmbeddingUrl === undefined) delete process.env.GITNEXUS_EMBEDDING_URL;
else process.env.GITNEXUS_EMBEDDING_URL = originalEmbeddingUrl;
if (originalEmbeddingModel === undefined) delete process.env.GITNEXUS_EMBEDDING_MODEL;
@ -190,7 +211,7 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => {
await run();
expect(acquireIndexLockMock).toHaveBeenCalledWith(dir);
expect(acquireIndexLockMock).toHaveBeenCalledWith(dir, { sweep: false });
expect(order[0]).toBe('lock');
expect(order.indexOf('loadMeta')).toBeGreaterThan(order.indexOf('lock'));
expect(order.indexOf('init')).toBeGreaterThan(order.indexOf('loadMeta'));
@ -319,7 +340,9 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => {
{ 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 { dir, lbugPath, metaPath } = await store();
await useRealIndexLock();
const lockPath = path.join(dir, 'analyze.lock');
const sourcePath = path.join(dir, 'lbug.staging.12345678-1234-4123-8123-123456789abc');
await writeFile(sourcePath, 'completed paid vectors');
await writeFile(`${sourcePath}.wal`, 'unfinished window');
@ -337,14 +360,17 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => {
},
});
await writeFile(metaPath, metadataBytes);
loadMetaMock.mockImplementation(async () => JSON.parse(await readFile(metaPath, 'utf8')));
loadMetaMock.mockImplementation(async () => {
expect(existsSync(lockPath)).toBe(true);
return 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(acquireIndexLockMock).toHaveBeenCalledWith(dir, { sweep: false });
expect(loadMetaMock.mock.invocationCallOrder[0]).toBeGreaterThan(
acquireIndexLockMock.mock.invocationCallOrder[0]!,
);
@ -355,11 +381,36 @@ describe('embeddingsSyncCommand writer safety (#3065)', () => {
expect(runEmbeddingPipelineMock).not.toHaveBeenCalled();
expect(saveMetaMock).not.toHaveBeenCalled();
expect(releaseMock).toHaveBeenCalledTimes(1);
expect(existsSync(lockPath)).toBe(false);
expect(await readFile(metaPath, 'utf8')).toBe(metadataBytes);
expect(await readFile(lbugPath, 'utf8')).toBe('db');
expect(await readFile(sourcePath, 'utf8')).toBe('completed paid vectors');
expect(await readFile(`${sourcePath}.wal`, 'utf8')).toBe('unfinished window');
});
it('sweeps orphaned staging files before writable sync work without a recovery receipt', async () => {
const { dir, metaPath } = await store();
await useRealIndexLock();
await writeFile(metaPath, JSON.stringify(BASE_META));
const sourcePath = path.join(dir, 'lbug.staging.12345678-1234-4123-8123-123456789abc');
await writeFile(sourcePath, 'orphaned database');
await writeFile(`${sourcePath}.wal`, 'orphaned WAL');
loadMetaMock.mockImplementation(async () => JSON.parse(await readFile(metaPath, 'utf8')));
ensurePrivateSharedGraphMock.mockImplementation(async () => {
expect(existsSync(path.join(dir, 'analyze.lock'))).toBe(true);
expect(existsSync(sourcePath)).toBe(false);
expect(existsSync(`${sourcePath}.wal`)).toBe(false);
return true;
});
await run();
expect(ensurePrivateSharedGraphMock).toHaveBeenCalledTimes(1);
expect(runEmbeddingPipelineMock).toHaveBeenCalledTimes(1);
expect(releaseMock).toHaveBeenCalledTimes(1);
expect(existsSync(path.join(dir, 'analyze.lock'))).toBe(false);
});
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

@ -205,9 +205,15 @@ describe('sweepStagingArtifacts', () => {
};
it('retains the checkpoint-referenced family while reclaiming unrelated orphans', () => {
const family = ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'].map(
(suffix) => recoveryStage + suffix,
);
const family = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
].map((suffix) => recoveryStage + suffix);
for (const name of [...family, `${recoveryStage}.unexpected`, 'lbug.staging.orphan']) {
writeFileSync(path.join(dir, name), 'x');
}

View file

@ -1,4 +1,5 @@
import fs from 'node:fs';
import { EventEmitter } from 'node:events';
import os from 'node:os';
import path from 'node:path';
import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest';
@ -10,8 +11,11 @@ const h = vi.hoisted(() => ({
connClose: vi.fn<() => Promise<void>>(),
query: vi.fn(),
abortBuilder: vi.fn(),
spawn: vi.fn(),
}));
vi.mock('node:child_process', () => ({ spawn: h.spawn }));
vi.mock('@ladybugdb/core', () => {
class Database {
constructor(...args: unknown[]) {
@ -44,6 +48,15 @@ vi.mock('../../src/core/embeddings/embedding-restore-spill.js', async (importOri
});
describe('staged embedding recovery child native lifecycle', () => {
const suffixes = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
];
let tmp: string;
let dbPath: string;
let exportDir: string;
@ -76,12 +89,24 @@ describe('staged embedding recovery child native lifecycle', () => {
});
afterEach(() => {
vi.useRealTimers();
process.argv = originalArgv;
process.exitCode = originalExitCode;
vi.restoreAllMocks();
fs.rmSync(tmp, { recursive: true, force: true });
});
function sourceFamily() {
return Object.fromEntries(
suffixes.map((suffix) => [suffix, fs.readFileSync(dbPath + suffix, 'utf8')]),
);
}
function seedCompleteFamily() {
for (const suffix of suffixes) fs.writeFileSync(dbPath + suffix, `retained ${suffix}`);
return sourceFamily();
}
async function runRejectedChild(message: string): Promise<void> {
await import('../../src/core/embeddings/staged-embedding-recovery-child.js');
await vi.waitFor(() => {
@ -136,6 +161,123 @@ describe('staged embedding recovery child native lifecycle', () => {
expect(h.abortBuilder).toHaveBeenCalledOnce();
});
it('confines writable replay and failed checkpoint close to a separate copied family', async () => {
const sourceBefore = seedCompleteFamily();
h.dbCtor.mockImplementation((openedPath: string) => {
for (const suffix of suffixes) {
expect(fs.readFileSync(openedPath + suffix, 'utf8')).toBe(sourceBefore[suffix]);
fs.writeFileSync(openedPath + suffix, `replayed ${suffix}`);
}
});
h.dbClose.mockImplementation(async () => {
const openedPath = h.dbCtor.mock.calls[0][0] as string;
fs.writeFileSync(openedPath, 'partial checkpoint');
throw new Error('checkpoint close failed');
});
await import('../../src/core/embeddings/staged-embedding-recovery-child.js');
await vi.waitFor(() => expect(process.exitCode).toBe(1));
expect(h.dbCtor.mock.calls[0][0]).not.toBe(dbPath);
expect(path.relative(exportDir, h.dbCtor.mock.calls[0][0] as string)).not.toMatch(/^\.\./);
expect(sourceFamily()).toEqual(sourceBefore);
expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false);
expect(process.stderr.write).toHaveBeenCalledWith('checkpoint close failed\n');
});
it('reclaims a timed-out writer copy without changing the retained source', async () => {
const sourceBefore = seedCompleteFamily();
h.query.mockReturnValue(new Promise(() => {}));
h.dbCtor.mockImplementation((openedPath: string) => {
for (const suffix of suffixes) fs.writeFileSync(openedPath + suffix, 'writer opened');
});
let childImport: Promise<unknown> | undefined;
const child = Object.assign(new EventEmitter(), {
stderr: new EventEmitter(),
kill: vi.fn(() => {
queueMicrotask(() => child.emit('close', null, 'SIGKILL'));
return true;
}),
});
h.spawn.mockImplementation((_command: string, args: string[]) => {
process.argv = [process.execPath, 'staged-embedding-recovery-child', ...args.slice(-3)];
childImport = import('../../src/core/embeddings/staged-embedding-recovery-child.js');
return child;
});
const { recoverStagedEmbeddings } =
await import('../../src/core/embeddings/staged-embedding-recovery.js');
vi.useFakeTimers();
const recovering = expect(
recoverStagedEmbeddings(dbPath, { dimensions: 2, timeoutMs: 500 }),
).rejects.toThrow(/timeout/);
await childImport;
expect(h.dbCtor).toHaveBeenCalledOnce();
const openedPath = h.dbCtor.mock.calls[0][0] as string;
await vi.advanceTimersByTimeAsync(500);
await recovering;
expect(child.kill).toHaveBeenCalledWith('SIGKILL');
expect(openedPath).not.toBe(dbPath);
expect(fs.existsSync(path.dirname(openedPath))).toBe(false);
expect(sourceFamily()).toEqual(sourceBefore);
expect(h.dbClose).not.toHaveBeenCalled();
});
it.each(['.wal', '.checkpoint.intent.lock', '.checkpoint.apply.lock'])(
'refuses a dangling family symlink before native open: %s',
async (suffix) => {
fs.rmSync(dbPath + suffix, { force: true });
fs.symlinkSync(path.join(tmp, 'missing-sidecar'), dbPath + suffix);
await import('../../src/core/embeddings/staged-embedding-recovery-child.js');
await vi.waitFor(() => expect(process.exitCode).toBe(1));
expect(h.dbCtor).not.toHaveBeenCalled();
expect(process.stderr.write).toHaveBeenCalledWith(
'staged embedding family is not a regular file\n',
);
expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false);
},
);
it('refuses a family entry replaced with a symlink between lstat and open', async () => {
const foreignPath = path.join(tmp, 'foreign-file');
fs.writeFileSync(foreignPath, 'foreign');
const open = fs.openSync;
const read = vi.spyOn(fs, 'readSync');
vi.spyOn(fs, 'openSync').mockImplementation((...args) => {
if (args[0] === dbPath) {
fs.rmSync(dbPath);
fs.symlinkSync(foreignPath, dbPath);
}
return open(...args);
});
await import('../../src/core/embeddings/staged-embedding-recovery-child.js');
await vi.waitFor(() => expect(process.exitCode).toBe(1));
expect(read).not.toHaveBeenCalled();
expect(h.dbCtor).not.toHaveBeenCalled();
expect(fs.readFileSync(foreignPath, 'utf8')).toBe('foreign');
expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false);
});
it('fails closed when copying the complete family runs out of space', async () => {
const sourceBefore = seedCompleteFamily();
vi.spyOn(fs, 'writeFileSync').mockImplementation(() => {
throw Object.assign(new Error('copy ran out of space'), { code: 'ENOSPC' });
});
await import('../../src/core/embeddings/staged-embedding-recovery-child.js');
await vi.waitFor(() => expect(process.exitCode).toBe(1));
expect(h.dbCtor).not.toHaveBeenCalled();
expect(sourceFamily()).toEqual(sourceBefore);
expect(fs.existsSync(path.join(exportDir, 'manifest.json'))).toBe(false);
expect(process.stderr.write).toHaveBeenCalledWith('copy ran out of space\n');
});
it('writes the manifest only after both native closes succeed', async () => {
const closed: string[] = [];
h.connClose.mockImplementation(async () => {