test(embeddings): cover recovery failure boundaries

This commit is contained in:
Gergő Magyar 2026-10-03 15:58:09 +01:00 • committed by GitHub
parent 26d9ba45be
commit c5e013a8d8
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
4 changed files with 710 additions and 7 deletions

View file

@ -0,0 +1,166 @@
import { afterEach, beforeEach, describe, expect, it } from 'vitest';
import { mkdtempSync, mkdirSync, rmSync, symlinkSync, writeFileSync } from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import {
readEmbeddingRecovery,
resolveEmbeddingRecovery,
} from '../../src/storage/embedding-recovery.js';
const stagingFile = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46';
const checkpoint = () => ({
kind: 'interrupted',
at: '2026-10-03T12:00:00.000Z',
nodesProcessed: 1,
totalNodes: 3,
chunksProcessed: 2,
model: 'test-model',
dimensions: 2,
provider: 'local',
pendingNodeIds: ['active'],
recovery: {
stagingFile,
schemaFingerprint: 'test-schema',
unsafeNodeIds: ['active', 'incomplete-restore'],
},
});
let dir: string;
beforeEach(() => {
dir = mkdtempSync(path.join(os.tmpdir(), 'gnx-embedding-recovery-'));
writeFileSync(path.join(dir, stagingFile), 'stage');
});
afterEach(() => rmSync(dir, { recursive: true, force: true }));
describe('resolveEmbeddingRecovery', () => {
it('resolves the exact staged generation and carries all unsafe node IDs', () => {
expect(resolveEmbeddingRecovery(dir, checkpoint())).toEqual({
...checkpoint().recovery,
dbPath: path.join(dir, stagingFile),
familyFiles: ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'].map(
(suffix) => stagingFile + suffix,
),
});
});
it.each([
'../lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46',
'/tmp/lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46',
'branches/foreign/lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46',
'..\\lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46',
'lbug',
'lbug.new',
'lbug.staging.orphan',
`${stagingFile}.wal`,
'lbug.staging.00000000-0000-4000-8000-000000000000',
])('rejects a foreign, unrecognized or missing generation: %s', (filename) => {
const marker = checkpoint();
marker.recovery.stagingFile = filename;
expect(resolveEmbeddingRecovery(dir, marker)).toBeUndefined();
});
it.each([
null,
[],
{ recovery: checkpoint().recovery },
{ ...checkpoint(), kind: 'partial' },
{ ...checkpoint(), kind: 'unverified-count' },
{ ...checkpoint(), kind: 'unknown' },
{ ...checkpoint(), at: 'invalid-time' },
{ ...checkpoint(), nodesProcessed: -1 },
{ ...checkpoint(), nodesProcessed: 4 },
{ ...checkpoint(), totalNodes: 1.5 },
{ ...checkpoint(), chunksProcessed: Number.NaN },
{ ...checkpoint(), model: '' },
{ ...checkpoint(), provider: null },
{ ...checkpoint(), dimensions: 0 },
{ ...checkpoint(), pendingNodeIds: [42] },
{ ...checkpoint(), pendingNodeIds: ['unexcluded'] },
{ ...checkpoint(), recovery: { ...checkpoint().recovery, schemaFingerprint: '' } },
{ ...checkpoint(), recovery: { ...checkpoint().recovery, unsafeNodeIds: undefined } },
{ ...checkpoint(), recovery: { ...checkpoint().recovery, unsafeNodeIds: [''] } },
])('rejects malformed or incompatible checkpoint shape %#', (marker) => {
expect(resolveEmbeddingRecovery(dir, marker)).toBeUndefined();
});
it('accepts a durable restored stage even before a new embedding window completes', () => {
expect(
resolveEmbeddingRecovery(dir, {
...checkpoint(),
nodesProcessed: 0,
chunksProcessed: 0,
pendingNodeIds: [],
}),
).toBeDefined();
});
it('rejects a directory in place of the staged database', () => {
rmSync(path.join(dir, stagingFile));
mkdirSync(path.join(dir, stagingFile));
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('rejects a dangling staged sidecar symlink', () => {
symlinkSync(path.join(dir, 'missing'), path.join(dir, `${stagingFile}.wal`));
expect(resolveEmbeddingRecovery(dir, checkpoint())).toBeUndefined();
});
});
describe('readEmbeddingRecovery', () => {
const writeMeta = (filename: string, marker: unknown = checkpoint()): void => {
writeFileSync(path.join(dir, filename), JSON.stringify({ embeddingCheckpoint: marker }));
};
it('prefers the primary metadata file over a stale legacy reference', () => {
writeMeta('gitnexus.json');
writeMeta('meta.json', null);
expect(readEmbeddingRecovery(dir)?.stagingFile).toBe(stagingFile);
});
it('loads the legacy mirror only when the primary metadata file is absent', () => {
writeMeta('meta.json');
expect(readEmbeddingRecovery(dir)?.stagingFile).toBe(stagingFile);
});
it('does not resurrect a legacy reference when the primary file is malformed', () => {
writeFileSync(path.join(dir, 'gitnexus.json'), '{');
writeMeta('meta.json');
expect(readEmbeddingRecovery(dir)).toBeUndefined();
});
it('does not fall back from valid primary metadata with no recovery reference', () => {
writeMeta('gitnexus.json', null);
writeMeta('meta.json');
expect(readEmbeddingRecovery(dir)).toBeUndefined();
});
it('rejects symlinked primary metadata rather than reading a foreign receipt', () => {
writeMeta('meta.json');
symlinkSync(path.join(dir, 'meta.json'), path.join(dir, 'gitnexus.json'));
expect(readEmbeddingRecovery(dir)).toBeUndefined();
});
it('resolves branch-slot provenance within that slot despite a flat storagePath', () => {
const branchSlot = path.join(dir, 'branches', 'feature');
mkdirSync(branchSlot, { recursive: true });
writeFileSync(path.join(branchSlot, stagingFile), 'branch-stage');
writeFileSync(
path.join(branchSlot, 'gitnexus.json'),
JSON.stringify({ storagePath: dir, embeddingCheckpoint: checkpoint() }),
);
expect(readEmbeddingRecovery(branchSlot)?.dbPath).toBe(path.join(branchSlot, stagingFile));
expect(readEmbeddingRecovery(dir)).toBeUndefined();
});
});

View file

@ -17,6 +17,7 @@ import {
existsSync,
chmodSync,
symlinkSync,
mkdirSync,
} from 'node:fs';
import os from 'node:os';
import path from 'node:path';
@ -178,6 +179,103 @@ describe('release', () => {
});
describe('sweepStagingArtifacts', () => {
const recoveryStage = 'lbug.staging.6e34c761-bf58-46cc-8b54-78d11607bc46';
const seedRecovery = (stagingFile = recoveryStage): void => {
writeFileSync(
path.join(dir, 'gitnexus.json'),
JSON.stringify({
embeddingCheckpoint: {
kind: 'interrupted',
at: '2026-10-03T12:00:00.000Z',
nodesProcessed: 1,
totalNodes: 2,
chunksProcessed: 1,
model: 'test-model',
dimensions: 2,
provider: 'local',
pendingNodeIds: ['node-2'],
recovery: {
stagingFile,
schemaFingerprint: 'schema-v1',
unsafeNodeIds: ['node-2'],
},
},
}),
);
};
it('retains the checkpoint-referenced family while reclaiming unrelated orphans', () => {
const family = ['', '.wal', '.shadow', '.wal.checkpoint', '.lock'].map(
(suffix) => recoveryStage + suffix,
);
for (const name of [...family, `${recoveryStage}.unexpected`, 'lbug.staging.orphan']) {
writeFileSync(path.join(dir, name), 'x');
}
seedRecovery();
sweepStagingArtifacts(dir);
for (const name of family) expect(existsSync(path.join(dir, name))).toBe(true);
expect(existsSync(path.join(dir, `${recoveryStage}.unexpected`))).toBe(false);
expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(false);
});
it('preserves a referenced generation when a shared caller acquires the lock', async () => {
writeFileSync(path.join(dir, recoveryStage), 'x');
seedRecovery();
const lock = await acquireIndexLock(dir);
try {
expect(existsSync(path.join(dir, recoveryStage))).toBe(true);
} finally {
lock.release();
}
});
it('retains only the referenced family in a branch sub-slot', async () => {
const branchDir = path.join(dir, 'branches', 'feature');
mkdirSync(branchDir, { recursive: true });
seedRecovery();
const metadata = JSON.parse(readFileSync(path.join(dir, 'gitnexus.json'), 'utf8'));
writeFileSync(
path.join(branchDir, 'gitnexus.json'),
JSON.stringify({ ...metadata, storagePath: dir }),
);
writeFileSync(path.join(branchDir, recoveryStage), 'stage');
writeFileSync(path.join(branchDir, 'lbug.staging.orphan'), 'orphan');
writeFileSync(path.join(dir, 'lbug.staging.orphan'), 'other-slot');
const lock = await acquireIndexLock(branchDir);
try {
expect(existsSync(path.join(branchDir, recoveryStage))).toBe(true);
expect(existsSync(path.join(branchDir, 'lbug.staging.orphan'))).toBe(false);
expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(true);
} finally {
lock.release();
}
});
it('rejects a symlinked reference and never follows it while reclaiming the stage', () => {
const external = path.join(dir, 'foreign-database');
writeFileSync(external, 'external');
symlinkSync(external, path.join(dir, recoveryStage));
seedRecovery();
sweepStagingArtifacts(dir);
expect(existsSync(path.join(dir, recoveryStage))).toBe(false);
expect(readFileSync(external, 'utf8')).toBe('external');
});
it('does not sweep any generation when acquisition explicitly defers cleanup', async () => {
writeFileSync(path.join(dir, 'lbug.staging.orphan'), 'x');
const lock = await acquireIndexLock(dir, { sweep: false });
try {
expect(existsSync(path.join(dir, 'lbug.staging.orphan'))).toBe(true);
} finally {
lock.release();
}
});
it('removes only staging files, never the live index or its sidecars', () => {
const files = [
'lbug',

View file

@ -7,7 +7,11 @@ import {
saveMeta,
type RepoMeta,
} from '../../src/storage/repo-manager.js';
import { EMBEDDING_DIMS, STALE_HASH_SENTINEL } from '../../src/core/lbug/schema.js';
import {
EMBEDDING_DIMS,
STALE_HASH_SENTINEL,
SCHEMA_FINGERPRINT,
} from '../../src/core/lbug/schema.js';
import { getIndexIncompleteReasons } from '../../src/core/index-freshness.js';
import type {
EmbeddingPipelineOptions,
@ -2375,19 +2379,27 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => {
});
expect(snapshots.windowStart?.lastCommit).not.toBe(currentCommit);
// ── Post-window: the one save that legitimately measured the count ──
// The measured count belongs to the unpublished staging generation.
expect(snapshots.postWindow).toMatchObject({
lastCommit: STALE_COMMIT,
fileHashes: STALE_HASHES,
incrementalInProgress: { phase: 'full-rebuild' },
stats: { embeddings: LIVE_EMBEDDING_COUNT },
stats: { embeddings: 7 },
});
expect(snapshots.postWindow?.embeddingCheckpoint?.recovery).toMatchObject({
stagingFile: expect.stringMatching(/^lbug\.staging\.[a-f0-9-]+$/),
schemaFingerprint: SCHEMA_FINGERPRINT,
unsafeNodeIds: [],
});
// ── Window 2: no stale restatement over the measured figure ────────
expect(snapshots.secondWindow).toMatchObject({
lastCommit: STALE_COMMIT,
stats: { embeddings: LIVE_EMBEDDING_COUNT },
embeddingCheckpoint: { pendingNodeIds: ['node-3', 'node-4'] },
stats: { embeddings: 7 },
embeddingCheckpoint: {
pendingNodeIds: ['node-3', 'node-4'],
recovery: { unsafeNodeIds: ['node-3', 'node-4'] },
},
});
// Only the finalize write — after the index is published — advances
@ -2461,7 +2473,7 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () =>
const mockResilienceHarness = (
controls: ResilienceControls,
): { runEmbeddingPipeline: Mock; loadCachedEmbeddings: Mock } => {
): { runEmbeddingPipeline: Mock; loadCachedEmbeddings: Mock; batchInsertEmbeddings: Mock } => {
const loadCachedEmbeddings = vi.fn(async () => ({
embeddingNodeIds: new Set<string>(),
embeddings: [],
@ -2530,11 +2542,13 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () =>
pipelineOptions: EmbeddingPipelineOptions,
): Promise<EmbeddingPipelineResult> => controls.pipeline(pipelineOptions),
);
const batchInsertEmbeddings = vi.fn(async () => undefined);
vi.doMock('../../src/core/embeddings/embedding-pipeline.js', () => ({
runEmbeddingPipeline,
batchInsertEmbeddings,
buildVectorIndex: vi.fn(async () => false),
}));
return { runEmbeddingPipeline, loadCachedEmbeddings };
return { runEmbeddingPipeline, loadCachedEmbeddings, batchInsertEmbeddings };
};
/** A checkpoint shaped exactly as `RepoMeta` declares it. */
@ -2589,6 +2603,8 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () =>
vi.doUnmock('../../src/storage/repo-manager.js');
vi.doUnmock('../../src/core/embeddings/embedding-identity.js');
vi.doUnmock('../../src/core/embeddings/embedding-pipeline.js');
vi.doUnmock('../../src/core/embeddings/staged-embedding-recovery.js');
vi.restoreAllMocks();
vi.resetModules();
vi.clearAllMocks();
vi.unstubAllEnvs();
@ -3039,4 +3055,332 @@ describe('runFullAnalysis embedding-checkpoint resilience (#2790 review)', () =>
await tmpRepo.cleanup();
}
});
const mockStagedFiles = async (): Promise<void> => {
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => {
if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture');
});
vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => {
await fs.rm(dbPath, { force: true });
});
};
const seedRecovery = async (storagePath: string, repoPath: string) => {
const stagingFile = 'lbug.staging.11111111-1111-4111-8111-111111111111';
const checkpoint = checkpointFixture({
kind: 'interrupted',
pendingNodeIds: [],
recovery: { stagingFile, schemaFingerprint: SCHEMA_FINGERPRINT, unsafeNodeIds: [] },
});
await seedMeta(storagePath, repoPath, {
stats: { nodes: 2, embeddings: 7 },
embeddingCheckpoint: checkpoint,
});
await fs.writeFile(`${storagePath}/${stagingFile}`, 'previous durable source');
vi.doMock('../../src/core/embeddings/staged-embedding-recovery.js', () => ({
recoverStagedEmbeddings: vi.fn(async () => ({ rows: [], embeddingNodeIds: new Set() })),
}));
return { checkpoint, stagingFile };
};
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-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } });
mockResilienceHarness({
count: [{ cnt: 9 }],
pipeline: async (options) => {
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 3,
chunksProcessed: 0,
nodeIds: ['active-node'],
});
expect((await loadMeta(storagePath))?.embeddingCheckpoint?.recovery).toBeUndefined();
await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 9 });
const midRun = await loadMeta(storagePath);
expect(midRun?.embeddingCheckpoint?.recovery).toBeUndefined();
expect(midRun?.stats?.embeddings).toBe(7);
return cleanResult();
},
});
expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toBeNull();
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined();
expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(9);
} finally {
await tmpRepo.cleanup();
}
});
it('keeps the previous durable source when manual checkpoints are disabled', async () => {
vi.stubEnv('GITNEXUS_WAL_MANUAL_CHECKPOINT', '0');
const tmpRepo = await createTempDir('gitnexus-3456-opt-out-source-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
const { checkpoint: original, stagingFile } = await seedRecovery(storagePath, tmpRepo.dbPath);
mockResilienceHarness({
count: [{ cnt: 9 }],
pipeline: async (options) => {
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 3,
chunksProcessed: 0,
nodeIds: ['active-node'],
});
await options.onCheckpoint?.({ nodesProcessed: 3, totalNodes: 3, chunksProcessed: 9 });
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original);
throw new Error('endpoint failure with manual checkpoints disabled');
},
});
await mockStagedFiles();
expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toMatchObject(
{ message: 'endpoint failure with manual checkpoints disabled' },
);
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original);
expect(await fs.readFile(`${storagePath}/${stagingFile}`, 'utf8')).toBe(
'previous durable source',
);
expect(
(await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')),
).toEqual([stagingFile]);
} finally {
await tmpRepo.cleanup();
}
});
it.skipIf(process.platform === 'win32')(
'retains the referenced stage through a symlinked storage directory',
async () => {
const tmpRepo = await createTempDir('gitnexus-3456-storage-alias-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
const actualStorage = `${tmpRepo.dbPath}/actual-index`;
await fs.mkdir(actualStorage);
await fs.symlink(actualStorage, storagePath, 'dir');
await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } });
mockResilienceHarness({
count: [{ cnt: 9 }],
pipeline: async (options) => {
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 3,
chunksProcessed: 0,
nodeIds: ['active-node'],
});
await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 3, chunksProcessed: 9 });
throw new Error('endpoint failed after durable window');
},
});
await mockStagedFiles();
expect(
await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, []),
).toMatchObject({ message: 'endpoint failed after durable window' });
const recovery = (await loadMeta(storagePath))?.embeddingCheckpoint?.recovery;
if (!recovery) throw new Error('expected retained recovery generation');
expect(await fs.readFile(`${actualStorage}/${recovery.stagingFile}`, 'utf8')).toBe(
'staged fixture',
);
expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(7);
} finally {
await tmpRepo.cleanup();
}
},
);
it.each(['checkpoint', 'metadata'] as const)(
'keeps the previous source when %s fails before recovery handoff',
async (failure) => {
const tmpRepo = await createTempDir('gitnexus-3456-handoff-failure-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
const { checkpoint: original, stagingFile } = await seedRecovery(
storagePath,
tmpRepo.dbPath,
);
const rename = fs.rename.bind(fs);
mockResilienceHarness({
count: [{ cnt: 9 }],
pipeline: async (options) => {
if (failure === 'metadata') {
vi.spyOn(fs, 'rename').mockImplementation(async (source, destination) => {
if (String(destination).endsWith('/gitnexus.json'))
throw new Error('metadata write failed');
return rename(source, destination);
});
} else {
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
vi.mocked(adapter.tryFlushWAL).mockResolvedValue(false);
}
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 3,
chunksProcessed: 0,
nodeIds: ['active-node'],
});
return cleanResult();
},
});
await mockStagedFiles();
const error = await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, []);
expect(error).toMatchObject({
message:
failure === 'metadata'
? 'metadata write failed'
: 'Could not checkpoint restored embeddings before recovery handoff.',
});
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toEqual(original);
expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(7);
expect(await fs.readFile(`${storagePath}/${stagingFile}`, 'utf8')).toBe(
'previous durable source',
);
expect(
(await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')),
).toEqual([stagingFile]);
} finally {
vi.restoreAllMocks();
await tmpRepo.cleanup();
}
},
);
it('keeps an active window unsafe when its completion checkpoint fails', async () => {
const tmpRepo = await createTempDir('gitnexus-3456-completion-failure-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
await seedMeta(storagePath, tmpRepo.dbPath, { stats: { embeddings: 7 } });
mockResilienceHarness({
count: [{ cnt: 9 }],
pipeline: async (options) => {
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 3,
chunksProcessed: 0,
nodeIds: ['active-node'],
});
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
vi.mocked(adapter.tryFlushWAL).mockResolvedValue(false);
await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 3, chunksProcessed: 9 });
return cleanResult();
},
});
await mockStagedFiles();
expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toMatchObject(
{ message: 'Could not checkpoint the completed embedding window for recovery.' },
);
const meta = await loadMeta(storagePath);
expect(meta?.stats?.embeddings).toBe(7);
expect(meta?.embeddingCheckpoint?.pendingNodeIds).toEqual(['active-node']);
expect(meta?.embeddingCheckpoint?.recovery?.unsafeNodeIds).toEqual(['active-node']);
if (!meta?.embeddingCheckpoint?.recovery) throw new Error('expected retained active window');
expect(
await fs.readFile(
`${storagePath}/${meta.embeddingCheckpoint.recovery.stagingFile}`,
'utf8',
),
).toBe('staged fixture');
} finally {
await tmpRepo.cleanup();
}
});
it('retains a failed stage and keeps future incomplete restore groups unsafe', async () => {
const tmpRepo = await createTempDir('gitnexus-3456-unsafe-restore-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
await seedMeta(storagePath, tmpRepo.dbPath, { stats: { nodes: 2, embeddings: 1 } });
let completedWindow: RepoMeta | null = null;
const { loadCachedEmbeddings, batchInsertEmbeddings } = mockResilienceHarness({
count: [{ cnt: 5 }],
pipeline: async (options) => {
await options.onCheckpointWindowStart?.({
nodesProcessed: 0,
totalNodes: 2,
chunksProcessed: 0,
nodeIds: ['earlier-window-node'],
});
await options.onCheckpoint?.({ nodesProcessed: 1, totalNodes: 2, chunksProcessed: 1 });
completedWindow = await loadMeta(storagePath);
throw new Error('Maximum database size exceeded');
},
});
loadCachedEmbeddings.mockResolvedValue({
embeddingNodeIds: new Set([RESILIENCE_NODE_ID]),
embeddings: [
{
nodeId: RESILIENCE_NODE_ID,
chunkIndex: 0,
startLine: 1,
endLine: 2,
contentHash: 'current-hash',
embedding: new Array(EMBEDDING_DIMS).fill(0),
},
],
});
batchInsertEmbeddings.mockRejectedValue(new Error('restore batch partially inserted'));
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => {
if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture');
});
vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => {
await fs.rm(dbPath, { force: true });
});
const logs: string[] = [];
expect(
await runAnalyze(
tmpRepo.dbPath,
{ embeddings: true, force: true, skipAgentsMd: true, skipSkills: true },
logs,
),
).toMatchObject({ message: 'Maximum database size exceeded' });
expect(completedWindow?.embeddingCheckpoint?.recovery?.unsafeNodeIds).toEqual([
RESILIENCE_NODE_ID,
]);
expect(completedWindow?.stats?.embeddings).toBe(1);
const recovery = (await loadMeta(storagePath))?.embeddingCheckpoint?.recovery;
if (!recovery) throw new Error('expected retained recovery generation');
expect(await fs.readFile(`${storagePath}/${recovery.stagingFile}`, 'utf8')).toBe(
'staged fixture',
);
expect(logs).toContainEqual(expect.stringContaining('GITNEXUS_LBUG_MAX_DB_SIZE'));
} finally {
await tmpRepo.cleanup();
}
});
it('reclaims a failed current stage before any recovery checkpoint exists', async () => {
const tmpRepo = await createTempDir('gitnexus-3456-no-checkpoint-');
try {
const { storagePath } = getStoragePaths(tmpRepo.dbPath);
await seedMeta(storagePath, tmpRepo.dbPath, {});
mockResilienceHarness({
count: [{ cnt: 0 }],
pipeline: async () => {
throw new Error('failed before first window');
},
});
const adapter = await import('../../src/core/lbug/lbug-adapter.js');
vi.mocked(adapter.initLbug).mockImplementation(async (dbPath) => {
if (dbPath.includes('.staging.')) await fs.writeFile(dbPath, 'staged fixture');
});
vi.mocked(adapter.wipeLbugDbFiles).mockImplementation(async (dbPath) => {
await fs.rm(dbPath, { force: true });
});
expect(
await runAnalyze(
tmpRepo.dbPath,
{ embeddings: true, force: true, skipAgentsMd: true, skipSkills: true },
[],
),
).toMatchObject({ message: 'failed before first window' });
expect(
(await fs.readdir(storagePath)).filter((name) => name.startsWith('lbug.staging.')),
).toEqual([]);
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined();
} finally {
await tmpRepo.cleanup();
}
});
});

View file

@ -0,0 +1,95 @@
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import { afterEach, describe, expect, it } from 'vitest';
import {
createCachedEmbeddingsBuilder,
disposeEmbeddingSpill,
finalizeCachedEmbeddingsSnapshot,
ingestCachedEmbeddingRow,
materializeCachedEmbeddings,
type CachedEmbeddingsSnapshot,
} from '../../src/core/embeddings/embedding-restore-spill.js';
import {
mergeRecoveredEmbeddings,
validateRecoveredNodeGroups,
} from '../../src/core/embeddings/staged-embedding-recovery.js';
describe('staged embedding recovery', () => {
const snapshots: CachedEmbeddingsSnapshot[] = [];
let tmp: string | undefined;
afterEach(() => {
for (const snapshot of snapshots) disposeEmbeddingSpill(snapshot.spill);
snapshots.length = 0;
if (tmp) fs.rmSync(tmp, { recursive: true, force: true });
tmp = undefined;
});
function snapshot(
rows: { nodeId: string; chunkIndex: number; contentHash?: string; embedding?: number[] }[],
) {
tmp ??= fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-test-'));
const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0, spillDir: tmp });
for (const row of rows) {
ingestCachedEmbeddingRow(
builder,
{ startLine: 1, endLine: 2, embedding: [1, 2], ...row },
true,
);
}
const result = finalizeCachedEmbeddingsSnapshot(builder);
snapshots.push(result);
return result;
}
it('accepts only complete groups with one content hash and unique contiguous chunk ordinals', () => {
const cached = snapshot([
{ nodeId: 'complete', chunkIndex: 1, contentHash: 'same' },
{ nodeId: 'complete', chunkIndex: 0, contentHash: 'same' },
{ nodeId: 'gap', chunkIndex: 1, contentHash: 'same' },
{ nodeId: 'duplicate', chunkIndex: 0, contentHash: 'same' },
{ nodeId: 'duplicate', chunkIndex: 0, contentHash: 'same' },
{ nodeId: 'mixed', chunkIndex: 0, contentHash: 'old' },
{ nodeId: 'mixed', chunkIndex: 1, contentHash: 'new' },
{ nodeId: 'no-hash', chunkIndex: 0 },
{ nodeId: 'unsafe', chunkIndex: 0, contentHash: 'same' },
]);
expect([...validateRecoveredNodeGroups(cached.rows, new Set(['unsafe']))]).toEqual([
'complete',
]);
});
it('replaces an entire published node group and keeps unrelated rows without retaining vector arrays', () => {
const live = snapshot([
{ nodeId: 'changed', chunkIndex: 0, contentHash: 'old' },
{ nodeId: 'changed', chunkIndex: 1, contentHash: 'old' },
{ nodeId: 'other', chunkIndex: 0, contentHash: 'other' },
]);
const recovered = snapshot([
{ nodeId: 'changed', chunkIndex: 0, contentHash: 'new', embedding: [7, 8] },
]);
const merged = mergeRecoveredEmbeddings(live, recovered);
snapshots.push(merged);
expect(merged.embeddings).toEqual([]);
expect(merged.rows).toHaveLength(2);
expect(materializeCachedEmbeddings(merged, merged.rows)).toEqual([
expect.objectContaining({ nodeId: 'other', contentHash: 'other' }),
expect.objectContaining({
nodeId: 'changed',
chunkIndex: 0,
contentHash: 'new',
embedding: [7, 8],
}),
]);
expect(fs.existsSync(live.spill.path)).toBe(true);
expect(fs.existsSync(recovered.spill.path)).toBe(true);
});
it('does not accept missing vector bytes during a merge', () => {
const cached = snapshot([{ nodeId: 'complete', chunkIndex: 0, contentHash: 'same' }]);
fs.truncateSync(cached.spill.path, 12);
expect(() => mergeRecoveredEmbeddings(snapshot([]), cached)).toThrow(
/short embedding spill read/,
);
});
});