fix(embeddings): reuse completed vectors after interrupted analyze (#3463)

This commit is contained in:
Gergő Magyar 2026-10-04 09:26:02 +01:00 • committed by GitHub
parent 5f9f95f224
commit 16d7e9477b
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
20 changed files with 3339 additions and 43 deletions

View file

@ -331,8 +331,8 @@ const respondOk = (_req, res) => {
//
// upstream request handler, replaceable mid-test via `ctx.handler`;
// null points the proxy at a port nothing ever listens on
// listenAfterMs bind the upstream this late, so the first attempt(s) hit
// ECONNREFUSED (a single-instance restart window)
// listenAfterMs bind the upstream this late after the first refused attempt
// (a single-instance restart window)
// schemeless drop http:// from GITNEXUS_UPSTREAM_URL, the way Render's
// `fromService: { property: hostport }` yields it
// env extra environment for docker-server.mjs
@ -366,18 +366,15 @@ async function withProxy(
})
: null;
// A late (or never) bind needs its port reserved up front; otherwise let the
// OS assign one at listen time.
const upstreamPort =
server && listenAfterMs === 0
? await new Promise((r) => server.listen(0, '127.0.0.1', () => r(server.address().port)))
: await getFreePort();
const bindTimer =
server && listenAfterMs > 0
? setTimeout(() => server.listen(upstreamPort, '127.0.0.1'), listenAfterMs)
: null;
// Keep the upstream port bound until the proxy port is chosen. Releasing it
// sooner lets the OS assign both services the same port and proxy to itself.
const reservation = server ?? createServer();
const upstreamPort = await new Promise((r) =>
reservation.listen(0, '127.0.0.1', () => r(reservation.address().port)),
);
const port = await getFreePort();
let bindTimer = null;
const target = `127.0.0.1:${upstreamPort}`;
const proc = spawnServerWithEnv(dir, port, {
GITNEXUS_UPSTREAM_URL: schemeless ? target : `http://${target}`,
@ -390,18 +387,25 @@ async function withProxy(
...env,
});
proc.stderr.setEncoding('utf8');
proc.stderr.on('data', (chunk) => {
const collectStderr = (chunk) => {
ctx.stderr += chunk;
});
// Process startup must not consume the restart window or skip the retry.
if (server && listenAfterMs > 0 && !bindTimer && ctx.stderr.includes('ECONNREFUSED; retry')) {
bindTimer = setTimeout(() => server.listen(upstreamPort, '127.0.0.1'), listenAfterMs);
}
};
proc.stderr.on('data', collectStderr);
try {
await waitForServer(port);
if (!server || listenAfterMs > 0) await new Promise((r) => reservation.close(r));
await fn(port, ctx);
} finally {
proc.stderr.off('data', collectStderr);
if (bindTimer) clearTimeout(bindTimer);
await killAndWait(proc);
if (server?.listening) {
server.closeAllConnections?.();
await new Promise((r) => server.close(r));
if (reservation.listening) {
reservation.closeAllConnections?.();
await new Promise((r) => reservation.close(r));
}
await rm(dir, { recursive: true, force: true });
}
@ -661,8 +665,8 @@ it('returns 502 when the upstream is unreachable', async () => {
// -- Connection-retry across an upstream restart window ---------------------
//
// `listenAfterMs: 400` binds the upstream late, so the first attempt hits
// ECONNREFUSED and must be retried — a single-instance restart. The default 3
// `listenAfterMs: 400` binds the upstream 400ms after the first ECONNREFUSED,
// so the request must be retried — a single-instance restart. The default 3
// attempts (backoff 250ms, 500ms) span ~750ms, so a retry lands after the bind.
it('retries a connection-refused POST and succeeds once the upstream is up', async () => {
@ -676,6 +680,7 @@ it('retries a connection-refused POST and succeeds once the upstream is up', asy
assert.match(res.body, /"ok":true/);
assert.equal(ctx.calls, 1, 'upstream must run the job exactly once (no double-execute)');
assert.equal(ctx.body, '{"repo":"x"}', 'buffered body replayed intact');
assert.match(ctx.stderr, /ECONNREFUSED; retry/, 'the restart gap must exercise a retry');
});
});

View file

@ -4,7 +4,11 @@ import { LBUG_DIRECTORY } from '../storage/storage-constants.js';
import path from 'node:path';
import { cliInfo } from './cli-message.js';
import { getGitRoot } from '../storage/git.js';
import { acquireIndexLock, requireExclusiveIndexLock } from '../storage/index-lock.js';
import {
acquireIndexLock,
requireExclusiveIndexLock,
sweepStagingArtifacts,
} from '../storage/index-lock.js';
import { getStoragePaths, loadMeta, saveMeta } from '../storage/repo-manager.js';
import {
closeLbug,
@ -50,12 +54,22 @@ export const embeddingsSyncCommand = async (inputPath?: string): Promise<void> =
// Writes go to the slot's own graph. A shared-store checkout that reads an
// immutable commit graph (#3352) takes a private copy first.
const lbugPath = path.join(metaDir, LBUG_DIRECTORY);
const lock = await acquireIndexLock(metaDir);
const lock = await acquireIndexLock(metaDir, { sweep: false });
try {
requireExclusiveIndexLock(
lock,
`Cannot acquire the index lock at ${metaDir}; refusing an unlocked embeddings sync.`,
);
// Sync writes the published graph and cannot recover a staged generation.
// Reject even malformed receipts before sweeping staging files or writing.
const recoveryCheckpoint = (await loadMeta(metaDir))?.embeddingCheckpoint;
if (recoveryCheckpoint && Object.hasOwn(recoveryCheckpoint, 'recovery')) {
throw new Error(
'Cannot sync embeddings: the index checkpoint references staged embeddings. ' +
'Run `gitnexus analyze` to recover them first.',
);
}
sweepStagingArtifacts(metaDir);
if (!(await ensurePrivateSharedGraph(metaDir, (m) => console.log(` ${m}`)))) {
throw new Error('The shared graph this checkout reads is gone. Run gitnexus analyze first.');
}

View file

@ -0,0 +1,181 @@
/** Isolated strict native reader. Never import this entrypoint into analyze. */
import fs from 'node:fs';
import path from 'node:path';
import lbug from '@ladybugdb/core';
import { createLbugDatabase, toNativeSafePath } from '../lbug/lbug-config.js';
import { FAMILY_SUFFIXES } from '../../storage/embedding-recovery.js';
import {
abortCachedEmbeddingsBuilder,
createCachedEmbeddingsBuilder,
finalizeCachedEmbeddingsSnapshot,
ingestCachedEmbeddingRow,
} from './embedding-restore-spill.js';
import type { StagedEmbeddingExport } from './staged-embedding-recovery.js';
/** Native replay and close may checkpoint; only give them disposable copies. */
function copyRecoveryFamily(dbPath: string, exportDir: string): string {
const replayDir = fs.mkdtempSync(path.join(exportDir, 'replay-'));
const replayPath = path.join(replayDir, path.basename(dbPath));
const noFollow = fs.constants.O_NOFOLLOW ?? 0;
const flags = fs.constants.O_RDONLY | noFollow | (fs.constants.O_NONBLOCK ?? 0);
const buffer = Buffer.allocUnsafe(1024 * 1024);
for (const suffix of FAMILY_SUFFIXES) {
const sourcePath = dbPath + suffix;
let entry: fs.BigIntStats;
try {
entry = fs.lstatSync(sourcePath, { bigint: true });
} catch (error) {
if (suffix && (error as NodeJS.ErrnoException).code === 'ENOENT') continue;
throw error;
}
if (!entry.isFile()) throw new Error('staged embedding family is not a regular file');
const source = fs.openSync(sourcePath, flags);
let destination: number | undefined;
try {
const opened = fs.fstatSync(source, { bigint: true });
const current = fs.lstatSync(sourcePath, { bigint: true });
if (
!opened.isFile() ||
!current.isFile() ||
(noFollow === 0 && opened.ino === 0n) ||
opened.dev !== entry.dev ||
opened.ino !== entry.ino ||
opened.dev !== current.dev ||
opened.ino !== current.ino
) {
throw new Error('staged embedding family changed while opening');
}
destination = fs.openSync(replayPath + suffix, 'wx', 0o600);
let copied = 0n;
for (;;) {
const bytesRead = fs.readSync(source, buffer, 0, buffer.length, null);
if (bytesRead === 0) break;
fs.writeFileSync(destination, buffer.subarray(0, bytesRead));
copied += BigInt(bytesRead);
}
const after = fs.fstatSync(source, { bigint: true });
if (
copied !== opened.size ||
after.size !== opened.size ||
after.mtimeNs !== opened.mtimeNs ||
after.ctimeNs !== opened.ctimeNs
) {
throw new Error('staged embedding family changed while copying');
}
} finally {
try {
if (destination !== undefined) fs.closeSync(destination);
} finally {
fs.closeSync(source);
}
}
}
return replayPath;
}
async function extract(): Promise<void> {
const [dbPath, exportDir, dimensionsArg] = process.argv.slice(2);
const dimensions = Number(dimensionsArg);
if (!dbPath || !exportDir || !Number.isInteger(dimensions) || dimensions <= 0) {
throw new Error('invalid staged embedding extraction arguments');
}
// The parent owns exportDir and reclaims it even after killing this child.
const replayPath = copyRecoveryFamily(dbPath, exportDir);
const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0, spillDir: exportDir });
const rejectedNodeIds = new Set<string>();
let db: lbug.Database | undefined;
let conn: lbug.Connection | undefined;
try {
// Avoid openLbugConnection's test-fixture lock sweep: a recovery source must
// never have its WAL removed, even when an external slot resembles a fixture.
db = createLbugDatabase(lbug, toNativeSafePath(replayPath), { throwOnWalReplayFailure: true });
conn = new lbug.Connection(db);
const queried = await conn.query(
'MATCH (e:CodeEmbedding) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.startLine AS startLine, e.endLine AS endLine, e.embedding AS embedding, e.contentHash AS contentHash',
);
const results = Array.isArray(queried) ? queried : [queried];
try {
if (results.length !== 1) throw new Error('unexpected staged embedding query result');
const result = results[0];
while (await result.hasNext()) {
const raw = await result.getNext();
const rec = raw as Record<string, unknown> & unknown[];
const nodeId = rec.nodeId ?? rec[0];
if (typeof nodeId !== 'string' || !nodeId)
throw new Error('invalid staged embedding node id');
const chunkIndex = rec.chunkIndex ?? rec[1];
const startLine = rec.startLine ?? rec[2];
const endLine = rec.endLine ?? rec[3];
const embedding = rec.embedding ?? rec[4];
const contentHash = rec.contentHash ?? rec[5];
const vector =
Array.isArray(embedding) ||
(ArrayBuffer.isView(embedding) && !(embedding instanceof DataView))
? Array.from(embedding as ArrayLike<number>)
: undefined;
if (
!Number.isInteger(chunkIndex) ||
Number(chunkIndex) < 0 ||
!Number.isInteger(startLine) ||
Number(startLine) < 0 ||
!Number.isInteger(endLine) ||
Number(endLine) < Number(startLine) ||
typeof contentHash !== 'string' ||
!contentHash ||
!vector ||
vector.length !== dimensions ||
vector.some(
(value) =>
typeof value !== 'number' ||
!Number.isFinite(value) ||
!Number.isFinite(Math.fround(value)),
)
) {
rejectedNodeIds.add(nodeId);
continue;
}
ingestCachedEmbeddingRow(
builder,
{ nodeId, chunkIndex, startLine, endLine, embedding: vector, contentHash },
true,
);
}
} finally {
for (const result of results) await result.close();
}
// Both closes must succeed. Suppressed native teardown errors are unsafe.
await conn.close();
await db.close();
const snapshot = finalizeCachedEmbeddingsSnapshot(builder);
if (snapshot.spill) fs.renameSync(snapshot.spill.path, path.join(exportDir, 'vectors.bin'));
const manifest: StagedEmbeddingExport = {
version: 1,
dimensions,
rows: snapshot.rows,
rejectedNodeIds: [...rejectedNodeIds],
};
fs.writeFileSync(path.join(exportDir, 'manifest.json'), JSON.stringify(manifest), {
flag: 'wx',
mode: 0o600,
});
} catch (err) {
abortCachedEmbeddingsBuilder(builder);
// Cleanup is best effort on a rejected source, never used to approve output.
try {
await conn?.close();
} catch {
/* rejected */
}
try {
await db?.close();
} catch {
/* rejected */
}
throw err;
}
}
extract().catch((err: unknown) => {
process.stderr.write(`${err instanceof Error ? err.message : String(err)}\n`);
process.exitCode = 1;
});

View file

@ -0,0 +1,244 @@
/** Recover paid embedding rows without opening an interrupted native DB in analyze. */
import { spawn } from 'node:child_process';
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
import {
abortCachedEmbeddingsBuilder,
createCachedEmbeddingsBuilder,
EmbeddingSpillReader,
emptyCachedEmbeddingsSnapshot,
finalizeCachedEmbeddingsSnapshot,
ingestCachedEmbeddingRow,
materializeCachedEmbeddings,
type CachedEmbeddingMeta,
type CachedEmbeddingsSnapshot,
} from './embedding-restore-spill.js';
export interface StagedEmbeddingRecoveryOptions {
dimensions: number;
/** Active checkpoint window and any incomplete inherited restore groups. */
excludedNodeIds?: Iterable<string>;
timeoutMs?: number;
}
export interface StagedEmbeddingExport {
version: 1;
dimensions: number;
rows: CachedEmbeddingMeta[];
rejectedNodeIds: string[];
}
const RESTORE_BATCH_SIZE = 200;
const DEFAULT_EXTRACTION_TIMEOUT_MS = 120_000;
/** A contiguous prefix is safe only when its node is outside every unsafe window. */
export function validateRecoveredNodeGroups(
rows: readonly CachedEmbeddingMeta[],
excludedNodeIds: ReadonlySet<string> = new Set(),
): Set<string> {
const groups = new Map<string, { hash: string; ordinals: Set<number>; invalid: boolean }>();
for (const row of rows) {
let group = groups.get(row.nodeId);
if (!group) {
group = { hash: row.contentHash, ordinals: new Set(), invalid: false };
groups.set(row.nodeId, group);
}
if (
typeof row.nodeId !== 'string' ||
!row.nodeId ||
typeof row.contentHash !== 'string' ||
!row.contentHash ||
row.contentHash !== group.hash ||
!Number.isInteger(row.chunkIndex) ||
row.chunkIndex < 0 ||
group.ordinals.has(row.chunkIndex) ||
!Number.isInteger(row.startLine) ||
row.startLine < 0 ||
!Number.isInteger(row.endLine) ||
row.endLine < row.startLine
)
group.invalid = true;
group.ordinals.add(row.chunkIndex);
}
const accepted = new Set<string>();
for (const [nodeId, group] of groups) {
if (group.invalid || excludedNodeIds.has(nodeId)) continue;
// Unique ordinals with max n-1 and zero present have no holes.
if (!group.ordinals.has(0)) continue;
if ([...group.ordinals].some((ordinal) => ordinal >= group.ordinals.size)) continue;
accepted.add(nodeId);
}
return accepted;
}
/** Whole recovered nodes replace whole published groups; vectors stay in bounded batches. */
export function mergeRecoveredEmbeddings(
live: CachedEmbeddingsSnapshot,
recovered: CachedEmbeddingsSnapshot,
): CachedEmbeddingsSnapshot {
const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0 });
try {
appendSnapshotRows(
live,
live.rows.filter((row) => !recovered.embeddingNodeIds.has(row.nodeId)),
builder,
);
appendSnapshotRows(recovered, recovered.rows, builder);
return finalizeCachedEmbeddingsSnapshot(builder);
} catch (err) {
abortCachedEmbeddingsBuilder(builder);
throw err;
}
}
function appendSnapshotRows(
snapshot: CachedEmbeddingsSnapshot,
rows: readonly CachedEmbeddingMeta[],
builder: ReturnType<typeof createCachedEmbeddingsBuilder>,
): void {
const reader = snapshot.spill ? new EmbeddingSpillReader(snapshot.spill) : undefined;
try {
for (let i = 0; i < rows.length; i += RESTORE_BATCH_SIZE) {
const batch = materializeCachedEmbeddings(
snapshot,
rows.slice(i, i + RESTORE_BATCH_SIZE),
reader,
);
for (const row of batch)
ingestCachedEmbeddingRow(builder, row as unknown as Record<string, unknown>, true);
}
} finally {
reader?.close();
}
}
/**
* Caller must validate checkpoint identity, schema, exact generation, and hold the
* index lock. A subprocess contains native WAL replay/query/destructor failures.
* A failed strict open is never retried with WAL removed or validation disabled.
*/
export async function recoverStagedEmbeddings(
dbPath: string,
options: StagedEmbeddingRecoveryOptions,
): Promise<CachedEmbeddingsSnapshot> {
if (!Number.isInteger(options.dimensions) || options.dimensions <= 0) {
throw new Error('invalid staged embedding dimensions');
}
const dbStat = fs.lstatSync(dbPath);
if (!dbStat.isFile() || dbStat.isSymbolicLink())
throw new Error('staged embedding DB is not a regular file');
const exportDir = fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-export-'));
try {
await runExtractionChild(dbPath, exportDir, options);
const manifestPath = path.join(exportDir, 'manifest.json');
const manifest = JSON.parse(fs.readFileSync(manifestPath, 'utf8')) as StagedEmbeddingExport;
if (
manifest.version !== 1 ||
manifest.dimensions !== options.dimensions ||
!Array.isArray(manifest.rows) ||
!Array.isArray(manifest.rejectedNodeIds) ||
manifest.rejectedNodeIds.some((id) => typeof id !== 'string') ||
manifest.rows.some((row, i) => !row || row.vectorIndex !== i)
)
throw new Error('invalid staged embedding export manifest');
if (manifest.rows.length === 0) return emptyCachedEmbeddingsSnapshot();
const spill = {
path: path.join(exportDir, 'vectors.bin'),
dims: options.dimensions,
rowCount: manifest.rows.length,
};
const vectorStat = fs.lstatSync(spill.path);
if (
!vectorStat.isFile() ||
vectorStat.isSymbolicLink() ||
vectorStat.size !== 12 + spill.rowCount * spill.dims * 4
) {
throw new Error('invalid staged embedding export size');
}
const excluded = new Set(options.excludedNodeIds ?? []);
for (const nodeId of manifest.rejectedNodeIds) excluded.add(nodeId);
const accepted = validateRecoveredNodeGroups(manifest.rows, excluded);
const exported: CachedEmbeddingsSnapshot = {
rows: manifest.rows,
embeddings: [],
embeddingNodeIds: accepted,
spill,
};
const builder = createCachedEmbeddingsBuilder({ inMemoryRowLimit: 0 });
const reader = new EmbeddingSpillReader(spill);
try {
const rows = manifest.rows.filter((row) => accepted.has(row.nodeId));
for (let i = 0; i < rows.length; i += RESTORE_BATCH_SIZE) {
for (const row of materializeCachedEmbeddings(
exported,
rows.slice(i, i + RESTORE_BATCH_SIZE),
reader,
)) {
if (
row.embedding.length !== options.dimensions ||
row.embedding.some((value) => !Number.isFinite(value))
) {
throw new Error('invalid staged embedding vector');
}
ingestCachedEmbeddingRow(builder, row as unknown as Record<string, unknown>, true);
}
}
return finalizeCachedEmbeddingsSnapshot(builder);
} catch (err) {
abortCachedEmbeddingsBuilder(builder);
throw err;
} finally {
reader.close();
}
} finally {
fs.rmSync(exportDir, { recursive: true, force: true });
}
}
function runExtractionChild(
dbPath: string,
exportDir: string,
options: StagedEmbeddingRecoveryOptions,
): Promise<void> {
const compiledPath = fileURLToPath(
new URL('./staged-embedding-recovery-child.js', import.meta.url),
);
const sourcePath = compiledPath.replace(/\.js$/, '.ts');
const childPath = fs.existsSync(compiledPath) ? compiledPath : sourcePath;
const args = childPath.endsWith('.ts')
? ['--import', import.meta.resolve('tsx'), childPath]
: [childPath];
args.push(dbPath, exportDir, String(options.dimensions));
return new Promise((resolve, reject) => {
const child = spawn(process.execPath, args, {
stdio: ['ignore', 'ignore', 'pipe'],
windowsHide: true,
});
let stderr = '';
let timedOut = false;
child.stderr.on('data', (chunk: Buffer) => {
if (stderr.length < 4096) stderr += chunk.toString().slice(0, 4096 - stderr.length);
});
const timer = setTimeout(() => {
timedOut = true;
// A process boundary is safe to kill even while native code is executing.
child.kill('SIGKILL');
}, options.timeoutMs ?? DEFAULT_EXTRACTION_TIMEOUT_MS);
child.once('error', (err) => {
clearTimeout(timer);
reject(err);
});
child.once('close', (code, signal) => {
clearTimeout(timer);
if (code === 0 && !signal && !timedOut) resolve();
else
reject(
new Error(
`staged embedding extraction failed${timedOut ? ' (timeout)' : ` (${signal ?? code})`}${stderr ? `: ${stderr.trim()}` : ''}`,
),
);
});
});
}

View file

@ -36,7 +36,12 @@ import fs from 'fs/promises';
import { constants as fsConstants, existsSync } from 'node:fs';
import { randomUUID } from 'node:crypto';
import { retryRename } from '../storage/fs-atomic.js';
import { acquireIndexLock, requireExclusiveIndexLock } from '../storage/index-lock.js';
import {
acquireIndexLock,
requireExclusiveIndexLock,
sweepStagingArtifacts,
} from '../storage/index-lock.js';
import { resolveEmbeddingRecovery } from '../storage/embedding-recovery.js';
import { invalidateNodeWorkspacePackages } from './ingestion/import-resolvers/node-workspace-packages.js';
import {
logNameFallbackSummary,
@ -1287,6 +1292,8 @@ export async function runFullAnalysis(
const log = (msg: string) => callbacks.onLog?.(stripControlCharacters(msg));
const acquireOpts = {
log,
// Resolve and validate the canonical slot under the lock before cleanup.
sweep: false,
onWaitStart: () =>
callbacks.onProgress('lock', 0, 'Waiting for another analyze to finish on this index…'),
};
@ -1345,6 +1352,7 @@ export async function runFullAnalysis(
}
const flatShared = writeTarget.placement.branch ? undefined : writeTarget.sharedStore;
if (flatShared) await seedSharedSlot(flatShared, repoPath, log);
sweepStagingArtifacts(writeTarget.metaDir, log);
const slotToLeave = options.noShare ? await optedInSlotToLeave(repoPath) : undefined;
const result = await runFullAnalysisInner(
repoPath,
@ -1842,7 +1850,13 @@ async function runFullAnalysisInner(
decision = decideEmbeddingResume(checkpoint, embeddingIdentityForRun, resumeOptions);
}
if (decision.action === 'abort') throw new Error(decision.error);
log(decision.log);
log(
decision.action === 'resume' && checkpoint.recovery
? `Previous analyze recorded an embedding checkpoint (${checkpoint.nodesProcessed}/` +
`${checkpoint.totalNodes} nodes); validating staged vectors before retrying ` +
`${decision.pendingNodeIds.size} pending node(s).`
: decision.log,
);
if (options.dropEmbeddings) {
// --drop-embeddings has always implied a rebuild here; the decision only
// covers the marker.
@ -2671,6 +2685,74 @@ async function runFullAnalysisInner(
}
}
// A checkpoint's pending decision and its paid, complete vectors are
// independent: --force discards the former but can still reuse the latter.
// Select only the explicitly referenced generation, never an orphan by age.
const stagedRecovery = resolveEmbeddingRecovery(metaDir, existingMeta?.embeddingCheckpoint);
const stagedCheckpoint = existingMeta?.embeddingCheckpoint;
const inheritedUnsafeNodeIds = new Set([
...pendingEmbeddingNodeIds,
...(stagedRecovery?.unsafeNodeIds ?? []),
]);
if (shouldLoadCache && !options.dropEmbeddings && stagedRecovery && stagedCheckpoint) {
if (!embeddingIdentityForRun) {
const { resolveEmbeddingIdentity } = await import('./embeddings/embedding-identity.js');
embeddingIdentityForRun = resolveEmbeddingIdentity();
}
const marker = stagedCheckpoint;
const matchesIdentity =
marker.model === embeddingIdentityForRun.model &&
marker.dimensions === embeddingIdentityForRun.dimensions &&
marker.provider === embeddingIdentityForRun.provider;
if (matchesIdentity && stagedRecovery.schemaFingerprint === SCHEMA_FINGERPRINT) {
// A force-discarded pending decision must not turn a known incomplete
// inherited group into a reusable cache merely because its hash matches.
if (inheritedUnsafeNodeIds.size > 0) {
const rows = cachedSnapshot.rows.filter((row) => !inheritedUnsafeNodeIds.has(row.nodeId));
cachedSnapshot = {
...cachedSnapshot,
rows,
embeddingNodeIds: new Set(rows.map((row) => row.nodeId)),
};
}
let recovered: CachedEmbeddingsSnapshot | undefined;
try {
const { recoverStagedEmbeddings, mergeRecoveredEmbeddings } =
await import('./embeddings/staged-embedding-recovery.js');
recovered = await recoverStagedEmbeddings(stagedRecovery.dbPath, {
dimensions: embeddingIdentityForRun.dimensions,
excludedNodeIds: stagedRecovery.unsafeNodeIds,
});
const liveCacheDims = snapshotEmbeddingDims(cachedSnapshot);
if (liveCacheDims !== undefined && liveCacheDims !== embeddingIdentityForRun.dimensions) {
log(
`Embedding dimensions changed (${liveCacheDims}d -> ` +
`${embeddingIdentityForRun.dimensions}d), discarding published cache`,
);
discardCachedEmbeddings();
}
if (recovered.rows.length > 0) {
const merged = mergeRecoveredEmbeddings(cachedSnapshot, recovered);
disposeEmbeddingSpill(cachedSnapshot.spill);
adoptCachedEmbeddings(merged);
}
log(
`Recovered ${recovered.rows.length} complete staged embedding chunk(s) ` +
`for ${recovered.embeddingNodeIds.size} node(s); unchanged content can reuse them.`,
);
} catch (err) {
log(
`Warning: could not recover staged embeddings (${(err as Error).message}); ` +
'the retry will regenerate missing chunks.',
);
} finally {
disposeEmbeddingSpill(recovered?.spill);
}
} else {
log('Staged embedding identity or schema changed; its vectors will not be reused.');
}
}
// ── Load incremental parse cache ──────────────────────────────────
// Content-addressed: `--force` reuses parser shards; `useParseCache: false`
// stages a new generation under a run-unique parse-rebuild.* dir and publishes
@ -4487,6 +4569,16 @@ async function runFullAnalysisInner(
embeddingIdentityForRun = resolveEmbeddingIdentity();
}
const embeddingIdentity = embeddingIdentityForRun;
const stagedRecoveryEnabled = useAtomicSwap && isManualCheckpointEnabled();
if (useAtomicSwap && !stagedRecoveryEnabled) {
log(
'Manual WAL checkpoints are disabled; new staged work cannot be recovered after interruption. ' +
'Any previous durable recovery source is retained until publication.',
);
}
const unsafeRecoveryNodeIds = new Set([...inheritedUnsafeNodeIds, ...restoreFailedNodeIds]);
let activeWindowNodeIds: string[] = [];
let recoveryGenerationDurable = false;
// Build a Map<nodeId, contentHash> from cached embeddings for incremental mode
let existingEmbeddings: Map<string, string> | undefined;
if (cachedSnapshot.embeddingNodeIds.size > 0) {
@ -4540,6 +4632,14 @@ async function runFullAnalysisInner(
stagedCheckpointEmbeddingCount = embeddings;
}
const latestMeta = (await loadMeta(metaDir)) ?? existingMeta;
// An in-place write or manual-checkpoint opt-out cannot create a new
// recoverable staged generation. Keep the complete previous receipt:
// updated progress or unsafe nodes would describe different source bytes.
const preservedRecoveryCheckpoint =
!stagedRecoveryEnabled &&
resolveEmbeddingRecovery(metaDir, latestMeta?.embeddingCheckpoint)
? latestMeta?.embeddingCheckpoint
: undefined;
// First-ever analyze of this repo: no meta exists on disk yet (the
// pre-wipe dirty stamp only fires when one does). Mint the minimum
// RepoMeta requires, with `lastCommit: ''` — never `currentCommit` —
@ -4550,6 +4650,11 @@ async function runFullAnalysisInner(
lastCommit: '',
indexedAt: new Date().toISOString(),
};
const interrupted = mintInterruptedCheckpoint(
embeddingIdentity,
checkpoint,
pendingNodeIds,
);
await saveMeta(metaDir, {
...base,
...(embeddings === undefined || buildPath !== lbugPath
@ -4557,11 +4662,18 @@ async function runFullAnalysisInner(
: { stats: { ...base.stats, embeddings } }),
// Written by a run that is still IN FLIGHT — see the `kind` doc in
// repo-manager.ts.
embeddingCheckpoint: mintInterruptedCheckpoint(
embeddingIdentity,
checkpoint,
pendingNodeIds,
),
embeddingCheckpoint: preservedRecoveryCheckpoint ?? {
...interrupted,
...(stagedRecoveryEnabled
? {
recovery: {
stagingFile: path.basename(buildPath),
schemaFingerprint: SCHEMA_FINGERPRINT,
unsafeNodeIds: [...unsafeRecoveryNodeIds],
},
}
: {}),
},
});
};
@ -4584,7 +4696,21 @@ async function runFullAnalysisInner(
{
forceReembedNodeIds: pendingEmbeddingNodeIds,
onCheckpointWindowStart: async ({ nodeIds, ...checkpoint }) => {
const handoff = stagedRecoveryEnabled && !recoveryGenerationDurable;
if (handoff) {
if (!(await checkpointOnce())) {
throw new Error(
'Could not checkpoint restored embeddings before recovery handoff.',
);
}
recoveryGenerationDurable = true;
}
activeWindowNodeIds = nodeIds;
for (const id of nodeIds) unsafeRecoveryNodeIds.add(id);
await saveEmbeddingCheckpoint(checkpoint, nodeIds);
// Reclaim the old source only after the new durable generation's
// reference is saved. Later windows retain this same generation.
if (handoff) sweepStagingArtifacts(metaDir, log);
},
// ── The mid-run count is a DIAGNOSTIC, not a gate (#2790) ──────
// This used to run the count query bare. THIS callback's rejection
@ -4598,7 +4724,12 @@ async function runFullAnalysisInner(
// touch stats.embeddings" signal — so the checkpoint still lands,
// with whatever count is already on disk left alone.
onCheckpoint: async (checkpoint) => {
await checkpointOnce();
const durable = await checkpointOnce();
if (stagedRecoveryEnabled && !durable) {
throw new Error('Could not checkpoint the completed embedding window for recovery.');
}
for (const id of activeWindowNodeIds) unsafeRecoveryNodeIds.delete(id);
activeWindowNodeIds = [];
const measured = await measurePersistedEmbeddingCount(executeQuery);
if (measured.kind === 'unknown') {
log(
@ -5085,6 +5216,7 @@ async function runFullAnalysisInner(
// is a crash-safety improvement: a failed swap leaves the previous index
// live and the next run recovers via the full-rebuild path.
await saveMeta(metaDir, meta);
sweepStagingArtifacts(metaDir, log);
// Registry freshness is published only after the graph and its metadata.
// A failed close, swap, or metadata save must leave the previous registry
@ -5294,7 +5426,13 @@ async function runFullAnalysisInner(
// rethrow below is the surface, and the lock's sweep remains the backstop.
if (useAtomicSwap && buildPath !== lbugPath) {
try {
await wipeLbugDbFiles(buildPath);
const recovery = resolveEmbeddingRecovery(
metaDir,
(await loadMeta(metaDir))?.embeddingCheckpoint,
);
// Both paths belong to this locked slot. The validated generation
// basename identifies the same file even through a directory alias.
if (recovery?.stagingFile !== path.basename(buildPath)) await wipeLbugDbFiles(buildPath);
} catch {
/* swallow — orphan reclamation must never mask the real failure */
}
@ -5306,6 +5444,12 @@ async function runFullAnalysisInner(
// IndexLockTimeoutError and other domain failures with `instanceof`.
recordLiveIndexMutationRisk(err);
}
if (/max(?:imum)?(?: database| db)? size|database size limit|maxDBSize/i.test(String(err))) {
log(
'The database size limit was reached. Set GITNEXUS_LBUG_MAX_DB_SIZE to a larger ' +
'byte limit before retrying analyze; retained complete embeddings can be reused.',
);
}
throw err;
}
}

View file

@ -12,6 +12,7 @@ import {
acquireIndexLock,
IndexLockTimeoutError,
requireExclusiveIndexLock,
sweepStagingArtifacts,
type IndexLockHandle,
} from '../storage/index-lock.js';
import { ensurePrivateSharedGraph } from '../core/shared-store-analyze.js';
@ -2152,11 +2153,21 @@ export const createServer = async (port: number, host: string = '127.0.0.1') =>
// for the whole embedding write, released in the finally below.
let slotLock: IndexLockHandle | undefined;
try {
slotLock = await acquireIndexLock(storagePath);
slotLock = await acquireIndexLock(storagePath, { sweep: false });
requireExclusiveIndexLock(
slotLock,
`Cannot acquire the index lock at ${storagePath}; refusing an unlocked embedding run.`,
);
// This writer cannot recover staged generations. Preserve their
// receipts, including malformed ones, before sweeping or writing.
const recoveryCheckpoint = (await loadMeta(storagePath))?.embeddingCheckpoint;
if (recoveryCheckpoint && Object.hasOwn(recoveryCheckpoint, 'recovery')) {
throw new Error(
'Cannot generate embeddings: the index checkpoint references staged embeddings. ' +
'Run `gitnexus analyze` to recover them first.',
);
}
sweepStagingArtifacts(storagePath);
// Writes go to the slot's own graph; a shared-store checkout
// reading an immutable commit graph (#3352) takes a private copy.
if (!(await ensurePrivateSharedGraph(storagePath, () => {}))) {

View file

@ -0,0 +1,167 @@
/**
* Filesystem-only staged embedding provenance. Keep this independent of native
* and model imports: every index-lock caller needs the retention decision.
*/
import {
closeSync,
constants,
fstatSync,
lstatSync,
openSync,
readFileSync,
realpathSync,
} from 'node:fs';
import path from 'node:path';
import type { EmbeddingRecoveryReference } from './repo-meta.js';
import { INDEX_METADATA_FILE, LEGACY_METADATA_FILE } from './storage-constants.js';
const STAGING_FILENAME =
/^lbug\.staging\.[0-9a-f]{8}-[0-9a-f]{4}-4[0-9a-f]{3}-[89ab][0-9a-f]{3}-[0-9a-f]{12}$/;
export const FAMILY_SUFFIXES = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
] as const;
export interface ResolvedEmbeddingRecovery extends EmbeddingRecoveryReference {
dbPath: string;
/** Exact basenames, never a prefix match that can preserve another generation. */
familyFiles: string[];
}
const isRecord = (value: unknown): value is Record<string, unknown> =>
value !== null && typeof value === 'object' && !Array.isArray(value);
const isNonemptyString = (value: unknown): value is string =>
typeof value === 'string' && value.length > 0;
const isNodeIds = (value: unknown): value is string[] =>
Array.isArray(value) && value.every(isNonemptyString);
const isCount = (value: unknown): value is number =>
typeof value === 'number' && Number.isSafeInteger(value) && value >= 0;
/**
* Resolve only an explicit interrupted-generation receipt in this canonical
* slot. This validates provenance, not native contents or current identity;
* those checks must pass separately before any rows can be reused.
*/
export const resolveEmbeddingRecovery = (
lockDir: string,
checkpoint: unknown,
): ResolvedEmbeddingRecovery | undefined => {
if (!isRecord(checkpoint) || !isRecord(checkpoint.recovery)) return undefined;
if (checkpoint.kind !== undefined && checkpoint.kind !== 'interrupted') return undefined;
if (
!isNonemptyString(checkpoint.at) ||
!Number.isFinite(Date.parse(checkpoint.at)) ||
!isCount(checkpoint.nodesProcessed) ||
!isCount(checkpoint.totalNodes) ||
checkpoint.nodesProcessed > checkpoint.totalNodes ||
!isCount(checkpoint.chunksProcessed) ||
!isNonemptyString(checkpoint.model) ||
!isCount(checkpoint.dimensions) ||
checkpoint.dimensions === 0 ||
!isNonemptyString(checkpoint.provider) ||
(checkpoint.pendingNodeIds !== undefined && !isNodeIds(checkpoint.pendingNodeIds))
) {
return undefined;
}
const recovery = checkpoint.recovery;
if (
!isNonemptyString(recovery.stagingFile) ||
!STAGING_FILENAME.test(recovery.stagingFile) ||
!isNonemptyString(recovery.schemaFingerprint) ||
!isNodeIds(recovery.unsafeNodeIds)
) {
return undefined;
}
const unsafeNodeIds = new Set(recovery.unsafeNodeIds);
if (
isNodeIds(checkpoint.pendingNodeIds) &&
checkpoint.pendingNodeIds.some((nodeId) => !unsafeNodeIds.has(nodeId))
) {
return undefined;
}
try {
const canonicalDir = realpathSync(lockDir);
if (!lstatSync(canonicalDir).isDirectory()) return undefined;
const familyFiles = FAMILY_SUFFIXES.map((suffix) => recovery.stagingFile + suffix);
for (const [index, filename] of familyFiles.entries()) {
try {
// lstat refuses both live and dangling symlinks, without following one
// to a database outside the slot. The base file must exist.
if (!lstatSync(path.join(canonicalDir, filename)).isFile()) return undefined;
} catch (error) {
if (index > 0 && (error as NodeJS.ErrnoException).code === 'ENOENT') continue;
return undefined;
}
}
return {
dbPath: path.join(canonicalDir, recovery.stagingFile),
stagingFile: recovery.stagingFile,
schemaFingerprint: recovery.schemaFingerprint,
unsafeNodeIds: [...unsafeNodeIds],
familyFiles,
};
} catch {
return undefined;
}
};
/** Synchronous mirror of loadMeta's primary-first, absent-only fallback rule. */
export const readEmbeddingRecovery = (lockDir: string): ResolvedEmbeddingRecovery | undefined => {
let metadataPath = path.join(lockDir, INDEX_METADATA_FILE);
let descriptor: number | undefined;
const noFollow = constants.O_NOFOLLOW ?? 0;
const flags = constants.O_RDONLY | noFollow | (constants.O_NONBLOCK ?? 0);
try {
try {
descriptor = openSync(metadataPath, flags);
} catch (error) {
const code = (error as NodeJS.ErrnoException).code;
if (code !== 'ENOENT' && code !== 'ENOTDIR') return undefined;
// Windows cannot open with O_NOFOLLOW: an open of a dangling symlink
// reports ENOENT, but that existing primary entry must prevent fallback.
try {
lstatSync(metadataPath);
return undefined;
} catch (statError) {
const statCode = (statError as NodeJS.ErrnoException).code;
if (statCode !== 'ENOENT' && statCode !== 'ENOTDIR') return undefined;
}
metadataPath = path.join(lockDir, LEGACY_METADATA_FILE);
descriptor = openSync(metadataPath, flags);
}
const opened = fstatSync(descriptor, { bigint: true });
const entry = lstatSync(metadataPath, { bigint: true });
// Check the opened file itself and match the current non-symlink entry.
// This also refuses replacement on platforms without O_NOFOLLOW.
if (
!opened.isFile() ||
!entry.isFile() ||
(noFollow === 0 && opened.ino === 0n) ||
opened.dev !== entry.dev ||
opened.ino !== entry.ino
) {
return undefined;
}
const meta: unknown = JSON.parse(readFileSync(descriptor, 'utf8'));
if (!isRecord(meta)) return undefined;
// storagePath describes the flat/cache root, including in branch-slot
// metadata. The current locked directory is the generation boundary.
return resolveEmbeddingRecovery(lockDir, meta.embeddingCheckpoint);
} catch {
return undefined;
} finally {
if (descriptor !== undefined) {
try {
closeSync(descriptor);
} catch {
/* best-effort */
}
}
}
};

View file

@ -62,6 +62,7 @@ import path from 'node:path';
import os from 'node:os';
import { randomBytes, randomUUID, createHash } from 'node:crypto';
import { isProcessAlive } from '../utils/process-identity.js';
import { readEmbeddingRecovery } from './embedding-recovery.js';
const LOCK_FILENAME = 'analyze.lock';
const LOCK_RECORD_VERSION = 1 as const;
@ -433,8 +434,9 @@ const deniedCreateHandle = (
/**
* Delete orphaned build/staging artifacts left in the lock directory by a
* crashed prior writer. Safe precisely because we hold the exclusive lock: no
* other writer can be creating these here right now, so anything present is a
* crash orphan. Matches this slot's staging files ONLY — never `lbug` itself,
* other writer can be creating these here right now. A checkpoint-referenced
* generation is retained for embedding recovery; all other stages are orphans.
* Matches this slot's staging files ONLY — never `lbug` itself,
* never `lbug.wal`/`lbug.shadow` (the LIVE index's own sidecars), and never a
* `branches/<slug>/` sub-slot (which owns its own lock + sweep). Non-recursive.
*/
@ -442,6 +444,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo
// Matches `lbug.new`, `lbug.new.wal`, `lbug.staging.<id>`, `lbug.staging.<id>.wal`, …
// Does NOT match `lbug`, `lbug.wal`, `lbug.shadow`.
const stagingRe = /^lbug\.(staging\..+|new(\..+)?)$/;
const retained = new Set(readEmbeddingRecovery(lockDir)?.familyFiles ?? []);
let removed = 0;
let entries: string[];
try {
@ -450,7 +453,7 @@ export const sweepStagingArtifacts = (lockDir: string, log?: (msg: string) => vo
return;
}
for (const name of entries) {
if (!stagingRe.test(name)) continue;
if (!stagingRe.test(name) || retained.has(name)) continue;
try {
unlinkSync(path.join(lockDir, name));
removed++;

View file

@ -43,6 +43,15 @@ export type ContentRetention = 'full' | 'symbol' | 'none';
export type FtsProfile = 'full' | 'symbol-no-file-content' | 'name-only';
export const CONTENT_RETENTION_SCHEMA_VERSION = 1;
/** Exact staged generation whose completed embedding groups can survive a retry. */
export interface EmbeddingRecoveryReference {
/** A run-minted basename within this metadata file's index slot. */
stagingFile: string;
schemaFingerprint: string;
/** Active-window and inherited incomplete groups: never reusable until completed. */
unsafeNodeIds: string[];
}
/**
* Versioned receipt for the analyzer process that produced an index.
*
@ -521,6 +530,11 @@ export interface RepoMeta {
* subset of their chunks; for `'partial'` they hold none.
*/
pendingNodeIds?: string[];
/**
* Interrupted atomic builds only. This does not publish the staged graph
* or advance live embedding statistics; it identifies a recovery source.
*/
recovery?: EmbeddingRecoveryReference;
};
/**
* Name of the git branch this index represents (#2106). Absent for the

View file

@ -0,0 +1,52 @@
import lbug from '@ladybugdb/core';
import fs from 'node:fs';
import { createLbugDatabase } from '../../../src/core/lbug/lbug-config.ts';
const [dbPath, mode] = process.argv.slice(2);
const db = createLbugDatabase(lbug, dbPath);
const conn = new lbug.Connection(db);
async function query(cypher) {
const queried = await conn.query(cypher);
for (const result of Array.isArray(queried) ? queried : [queried]) {
await result.getAll();
await result.close();
}
}
await query('CREATE NODE TABLE CodeEmbedding (id STRING, nodeId STRING, chunkIndex INT32, startLine INT64, endLine INT64, embedding FLOAT[2], contentHash STRING, PRIMARY KEY(id))');
async function row(id, nodeId, chunkIndex, hash = 'same', startLine = 1, endLine = 3) {
await query(`CREATE (:CodeEmbedding {id: '${id}', nodeId: '${nodeId}', chunkIndex: ${chunkIndex}, startLine: ${startLine}, endLine: ${endLine}, embedding: [1.0, 2.0], contentHash: ${hash === null ? 'NULL' : `'${hash}'`}})`);
}
await row('complete-0', 'complete', 0);
await row('complete-1', 'complete', 1);
await row('other', 'other', 0);
await query('CHECKPOINT');
if (mode === 'interrupted-checkpoint') {
await row('checkpoint-only', 'checkpoint-only', 0);
const main = fs.readFileSync(dbPath);
const wal = fs.readFileSync(`${dbPath}.wal`);
await conn.close();
await db.close();
// Restore the pre-close bytes: writable close has already checkpointed them.
fs.writeFileSync(dbPath, main);
fs.writeFileSync(`${dbPath}.wal.checkpoint`, wal);
fs.writeFileSync(`${dbPath}.wal`, '');
fs.writeFileSync(`${dbPath}.shadow`, '');
fs.writeFileSync(`${dbPath}.checkpoint.intent.lock`, '');
fs.writeFileSync(`${dbPath}.checkpoint.apply.lock`, '');
} else if (mode === 'hard-kill') {
await row('unsafe', 'unsafe-prefix', 0);
process.kill(process.pid, 'SIGKILL');
} else {
await row('gap', 'gap', 1);
await row('duplicate-0', 'duplicate', 0);
await row('duplicate-1', 'duplicate', 0);
await row('mixed-0', 'mixed', 0, 'old');
await row('mixed-1', 'mixed', 1, 'new');
await row('missing-hash-0', 'missing-hash', 0);
await row('missing-hash-1', 'missing-hash', 1, null);
await row('bad-line', 'bad-line', 0, 'same', -1, 3);
await row('nan-0', 'nan', 0);
await query("CREATE (:CodeEmbedding {id: 'nan-1', nodeId: 'nan', chunkIndex: 1, startLine: 1, endLine: 3, embedding: [CAST('NaN', 'FLOAT'), 2.0], contentHash: 'same'})");
await conn.close();
await db.close();
}

View file

@ -0,0 +1,449 @@
/**
* Exercise the CLI -> native staged DB -> durable checkpoint -> SIGKILL ->
* isolated recovery -> graph rebuild -> publication chain. The endpoint is
* local and deterministic; request text is the billing/reuse oracle.
*/
import { spawn, spawnSync, type ChildProcess } from 'node:child_process';
import fs from 'node:fs';
import http from 'node:http';
import os from 'node:os';
import path from 'node:path';
import { fileURLToPath, pathToFileURL } from 'node:url';
import { afterAll, beforeAll, describe, expect, it } from 'vitest';
import { CLI_SPAWN_PREFIX, tsxLoaderUrl } from '../helpers/cli-entry.js';
const DIMS = 8;
const NODE_COUNT = 5_128;
const DEADLINE = process.env.CI ? 180_000 : 120_000;
const packageRoot = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..', '..');
const adapterUrl = pathToFileURL(path.join(packageRoot, 'src/core/lbug/lbug-adapter.ts')).href;
interface CheckpointMeta {
repoPath: string;
stats?: { embeddings?: number };
embeddingCheckpoint?: {
nodesProcessed: number;
chunksProcessed: number;
pendingNodeIds?: string[];
recovery?: { stagingFile: string; schemaFingerprint: string; unsafeNodeIds: string[] };
};
}
interface Run {
child: ChildProcess;
output: () => string;
done: Promise<{ code: number | null; signal: NodeJS.Signals | null; output: string }>;
}
let root: string;
let stoppedRepo: string;
let server: http.Server;
let endpoint: string;
let submitted: string[] = [];
let completed: string[] = [];
let durableTexts: string[] = [];
const activeWindowCompleted: string[] = [];
let held: string[] = [];
let stopAfterCheckpoint = false;
let currentRepo: string;
let gateResolve: (() => void) | undefined;
const running = new Set<ChildProcess>();
function readMeta(repo: string): CheckpointMeta {
return JSON.parse(fs.readFileSync(path.join(repo, '.gitnexus', 'gitnexus.json'), 'utf8'));
}
function cliEnv(repo: string): NodeJS.ProcessEnv {
const env = { ...process.env };
for (const key of Object.keys(env)) {
if (key.startsWith('GITNEXUS_EMBEDDING_') || key.startsWith('GITNEXUS_STORAGE_'))
delete env[key];
}
return {
...env,
GITNEXUS_HOME: path.join(root, `home-${path.basename(repo)}`),
GITNEXUS_SHARED_STORE: 'off',
GITNEXUS_LBUG_EXTENSION_INSTALL: 'never',
GITNEXUS_LBUG_BUFFER_POOL_SIZE: String(256 * 1024 * 1024),
GITNEXUS_EMBEDDING_URL: endpoint,
GITNEXUS_EMBEDDING_MODEL: 'staged-recovery-fixture',
GITNEXUS_EMBEDDING_DIMS: String(DIMS),
GITNEXUS_EMBEDDING_BATCH_SIZE: '5000',
GITNEXUS_EMBEDDING_SUB_BATCH_SIZE: '64',
GITNEXUS_EMBEDDING_MAX_ATTEMPTS: '1',
GITNEXUS_EMBEDDING_CACHE_IN_MEMORY_LIMIT: '0',
GITNEXUS_MEMORY: 'off',
// SIGKILL skips process exit hooks; keep orphaned cache/export spills in
// this suite's owned directory so afterAll can remove them as well.
TMPDIR: path.join(root, 'tmp'),
NODE_OPTIONS: `${process.env.NODE_OPTIONS || ''} --max-old-space-size=2048`.trim(),
CI: '1',
};
}
function runAnalyze(repo: string, flags: string[] = [], onOutput?: (output: string) => void): Run {
let output = '';
const child = spawn(
process.execPath,
[
...CLI_SPAWN_PREFIX,
'analyze',
repo,
'--no-share',
'--skip-skills',
'--skip-fts',
'--workers',
'1',
...flags,
],
{ cwd: repo, env: cliEnv(repo), stdio: ['ignore', 'pipe', 'pipe'] },
);
running.add(child);
const timeout = setTimeout(() => child.kill('SIGKILL'), DEADLINE);
const collect = (chunk: Buffer) => {
output += chunk.toString();
onOutput?.(output);
};
child.stdout?.on('data', collect);
child.stderr?.on('data', collect);
const done = new Promise<{ code: number | null; signal: NodeJS.Signals | null; output: string }>(
(resolve, reject) => {
child.once('error', reject);
child.once('close', (code, signal) => {
clearTimeout(timeout);
running.delete(child);
resolve({ code, signal, output });
});
},
);
return { child, done, output: () => output };
}
async function successfulAnalyze(repo: string, flags: string[] = []): Promise<string[]> {
submitted = [];
stopAfterCheckpoint = false;
currentRepo = repo;
const result = await runAnalyze(repo, flags).done;
expect(result.output, `CLI exited ${result.code}, signal ${result.signal}`).not.toContain(
'SIGABRT',
);
expect(result.code, result.output).toBe(0);
return [...submitted];
}
function cloneStoppedRepo(name: string): string {
const repo = path.join(root, name);
fs.cpSync(stoppedRepo, repo, { recursive: true });
for (const filename of ['gitnexus.json', 'meta.json']) {
const target = path.join(repo, '.gitnexus', filename);
const meta = JSON.parse(fs.readFileSync(target, 'utf8'));
meta.repoPath = repo;
meta.storagePath = path.join(repo, '.gitnexus');
fs.writeFileSync(target, JSON.stringify(meta));
}
return repo;
}
function readPublishedRows(
repo: string,
): Array<{ nodeId: string; chunkIndex: number; embedding: number[] }> {
// Keep native handles out of the vitest fork, and wait for clean teardown
// before accepting the receipt. Every read opens only a published DB.
const receiptPath = path.join(root, `rows-${path.basename(repo)}.json`);
const script = `
const adapter = await import(${JSON.stringify(adapterUrl)});
const fs = await import('node:fs');
await adapter.initLbug(${JSON.stringify(path.join(repo, '.gitnexus', 'lbug'))});
try {
const rows = await adapter.executeQuery('MATCH (e:CodeEmbedding) RETURN e.nodeId AS nodeId, e.chunkIndex AS chunkIndex, e.embedding AS embedding');
fs.writeFileSync(${JSON.stringify(receiptPath)}, JSON.stringify(rows));
console.log('ROWS_RECEIPT:' + rows.length);
} finally { await adapter.closeLbug(); }
`;
const result = spawnSync(
process.execPath,
['--import', tsxLoaderUrl(), '--input-type=module', '-e', script],
{
cwd: packageRoot,
env: cliEnv(repo),
encoding: 'utf8',
timeout: 30_000,
},
);
try {
const diagnostic =
`${result.error ?? ''} ${result.signal ?? ''}\n${result.stderr}\n${result.stdout}`.slice(
0,
4096,
);
expect(result.status, diagnostic).toBe(0);
expect(result.stdout).toMatch(/ROWS_RECEIPT:\d+/);
return JSON.parse(fs.readFileSync(receiptPath, 'utf8'));
} finally {
fs.rmSync(receiptPath, { force: true });
}
}
function expectPublishedComplete(repo: string, expectedNodes: number): void {
const meta = readMeta(repo);
expect(meta.embeddingCheckpoint).toBeUndefined();
const rows = readPublishedRows(repo);
expect(new Set(rows.map((row) => row.nodeId)).size).toBe(expectedNodes);
expect(meta.stats?.embeddings).toBe(rows.length);
const byNode = new Map<string, number[]>();
for (const row of rows) {
expect(row.embedding).toHaveLength(DIMS);
expect(row.embedding.every(Number.isFinite)).toBe(true);
const indices = byNode.get(row.nodeId) ?? [];
indices.push(row.chunkIndex);
byNode.set(row.nodeId, indices);
}
for (const indices of byNode.values()) {
expect(indices.sort((a, b) => a - b)).toEqual(
Array.from({ length: indices.length }, (_, index) => index),
);
}
expect(
fs.readdirSync(path.join(repo, '.gitnexus')).filter((name) => name.startsWith('lbug.staging.')),
).toEqual([]);
}
beforeAll(async () => {
if (process.platform === 'win32') return;
root = fs.mkdtempSync(path.join(os.tmpdir(), 'gn-staged-recovery-e2e-'));
fs.mkdirSync(path.join(root, 'tmp'));
stoppedRepo = path.join(root, 'interrupted');
currentRepo = stoppedRepo;
fs.mkdirSync(stoppedRepo);
const shortFunctions = Array.from(
{ length: NODE_COUNT - 1 },
(_, index) => `export function recoverable${index}() { return ${index}; }`,
);
// Multi-chunk nodes must be reused as a complete group, including their tail.
const longFunction = `export function longRecoverable() {\n${Array.from({ length: 80 }, (_, index) => ` // retained chunk marker ${index} ${'x'.repeat(80)}`).join('\n')}\n return 42;\n}`;
fs.writeFileSync(
path.join(stoppedRepo, 'functions.ts'),
`${longFunction}\n${shortFunctions.join('\n')}\n`,
);
const gitEnv = {
...process.env,
GIT_AUTHOR_NAME: 'test',
GIT_AUTHOR_EMAIL: 'test@test',
GIT_COMMITTER_NAME: 'test',
GIT_COMMITTER_EMAIL: 'test@test',
};
for (const args of [['init'], ['add', 'functions.ts'], ['commit', '-m', 'recovery fixture']]) {
const result = spawnSync('git', args, { cwd: stoppedRepo, env: gitEnv, encoding: 'utf8' });
expect(result.status, result.stderr).toBe(0);
}
server = http.createServer((request, response) => {
let body = '';
request.on('data', (chunk) => {
body += chunk;
});
request.on('end', () => {
const { input } = JSON.parse(body) as { input: string[] };
submitted.push(...input);
if (
stopAfterCheckpoint &&
(readMeta(currentRepo).embeddingCheckpoint?.nodesProcessed ?? 0) >= 5000
) {
if (activeWindowCompleted.length > 0) {
held = [...input];
gateResolve?.();
// One sub-batch has already inserted rows inside the unsafe window.
// Awaiting the next real request gates a crash before it completes.
return;
}
durableTexts = [...completed];
activeWindowCompleted.push(...input);
}
completed.push(...input);
response.writeHead(200, { 'content-type': 'application/json' });
response.end(
JSON.stringify({
data: input.map((text, index) => ({
index,
embedding: Array.from(
{ length: DIMS },
(_, dimension) => ((text.length + dimension * 17) % 101) / 101,
),
})),
}),
);
});
});
await new Promise<void>((resolve) => server.listen(0, '127.0.0.1', resolve));
const address = server.address() as { port: number };
endpoint = `http://127.0.0.1:${address.port}/v1`;
await successfulAnalyze(stoppedRepo, ['--index-only']);
expect(readMeta(stoppedRepo).stats?.embeddings ?? 0).toBe(0);
completed = [];
held = [];
stopAfterCheckpoint = true;
const gate = new Promise<void>((resolve) => {
gateResolve = resolve;
});
const first = runAnalyze(stoppedRepo, ['--force', '--embeddings']);
await Promise.race([
gate,
first.done.then((result) => {
throw new Error(`CLI exited before recovery gate: ${result.output}`);
}),
]);
first.child.kill('SIGKILL');
expect((await first.done).signal).toBe('SIGKILL');
const interrupted = readMeta(stoppedRepo);
expect(interrupted.embeddingCheckpoint?.nodesProcessed).toBe(5000);
expect(interrupted.embeddingCheckpoint?.recovery?.stagingFile).toMatch(/^lbug\.staging\./);
expect(interrupted.stats?.embeddings ?? 0).toBe(0);
expect(completed.length).toBeGreaterThanOrEqual(5000);
expect(completed.filter((text) => text.includes('longRecoverable')).length).toBeGreaterThan(1);
expect(activeWindowCompleted).toHaveLength(64);
expect(held.length).toBeGreaterThan(0);
gateResolve = undefined;
stopAfterCheckpoint = false;
}, DEADLINE * 2);
afterAll(async () => {
for (const child of running) child.kill('SIGKILL');
await Promise.all(
[...running].map(
(child) => new Promise<void>((resolve) => child.once('close', () => resolve())),
),
);
if (server) {
server.closeAllConnections();
await new Promise<void>((resolve) => server.close(() => resolve()));
}
if (root) fs.rmSync(root, { recursive: true, force: true });
});
// Atomic publication is POSIX-specific; Windows uses the in-place path.
describe
.skipIf(process.platform === 'win32')
.sequential('interrupted staged embedding recovery (real CLI and native DB)', () => {
it.each([
['plain', []],
['forced', ['--force', '--embeddings']],
] as const)(
'%s retry bills only unfinished groups and publishes an honest count',
async (name, flags) => {
const repo = cloneStoppedRepo(name);
const texts = await successfulAnalyze(repo, [...flags]);
const durable = new Set(durableTexts);
expect(texts.filter((text) => durable.has(text))).toEqual([]);
for (const text of held) expect(texts).toContain(text);
for (const text of activeWindowCompleted) expect(texts).toContain(text);
expect(texts.length).toBe(128);
expectPublishedComplete(repo, NODE_COUNT);
},
DEADLINE,
);
it(
'manual checkpoint opt-out completes a staged retry without resubmitting durable chunks',
async () => {
const repo = cloneStoppedRepo('manual-checkpoint-opt-out');
const previous = process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT;
process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT = '0';
try {
const texts = await successfulAnalyze(repo, ['--force', '--embeddings']);
expect(texts.filter((text) => new Set(durableTexts).has(text))).toEqual([]);
expect(texts.length).toBe(128);
expectPublishedComplete(repo, NODE_COUNT);
} finally {
if (previous === undefined) delete process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT;
else process.env.GITNEXUS_WAL_MANUAL_CHECKPOINT = previous;
}
},
DEADLINE,
);
it(
'survives another crash after harvesting but before the replacement is durable',
async () => {
const repo = cloneStoppedRepo('crash-again');
const source = readMeta(repo).embeddingCheckpoint?.recovery?.stagingFile;
expect(source).toBeDefined();
if (!source) throw new Error('fixture has no retained generation');
submitted = [];
currentRepo = repo;
let killed = false;
const retry = runAnalyze(repo, [], (output) => {
if (!killed && /Recovered \d+ complete staged embedding chunk/.test(output)) {
killed = true;
retry.child.kill('SIGKILL');
}
});
const result = await retry.done;
expect(killed, result.output).toBe(true);
expect(result.signal).toBe('SIGKILL');
expect(submitted).toEqual([]);
expect(readMeta(repo).embeddingCheckpoint?.recovery?.stagingFile).toBe(source);
expect(fs.existsSync(path.join(repo, '.gitnexus', source))).toBe(true);
const finalTexts = await successfulAnalyze(repo);
expect(finalTexts.filter((text) => new Set(durableTexts).has(text))).toEqual([]);
expect(finalTexts.length).toBe(128);
expectPublishedComplete(repo, NODE_COUNT);
},
DEADLINE * 2,
);
it(
'regenerates changed content and removes deleted nodes while reusing the other complete groups',
async () => {
const repo = cloneStoppedRepo('changed');
const source = path.join(repo, 'functions.ts');
const content = fs
.readFileSync(source, 'utf8')
.replace(
'export function recoverable5() { return 5; }',
'export function recoverable5() { return 999999; }',
)
.replace(
'export function recoverable6() { return 6; }',
'// deleted function retains line offsets',
);
fs.writeFileSync(source, content);
const texts = await successfulAnalyze(repo);
expect(texts.some((text) => text.includes('recoverable5') && text.includes('999999'))).toBe(
true,
);
expect(texts.some((text) => text.includes('recoverable6'))).toBe(false);
expect(texts.filter((text) => new Set(durableTexts).has(text))).toEqual([]);
expect(texts.length).toBe(129);
expectPublishedComplete(repo, NODE_COUNT - 1);
},
DEADLINE,
);
it(
'a forced retry with a different model does not import the staged cache',
async () => {
const repo = cloneStoppedRepo('different-model');
const texts = await successfulAnalyze(repo, [
'--force',
'--embeddings',
'--embedding-model',
'different-model',
]);
for (const text of durableTexts) expect(texts).toContain(text);
expect(texts.length).toBeGreaterThanOrEqual(NODE_COUNT);
expectPublishedComplete(repo, NODE_COUNT);
},
DEADLINE,
);
it(
'explicit drop abandons staged vectors without contacting the provider',
async () => {
const repo = cloneStoppedRepo('drop');
expect(await successfulAnalyze(repo, ['--force', '--drop-embeddings'])).toEqual([]);
expect(readMeta(repo).embeddingCheckpoint).toBeUndefined();
expect(readMeta(repo).stats?.embeddings ?? 0).toBe(0);
expect(readPublishedRows(repo)).toEqual([]);
},
DEADLINE,
);
});

View file

@ -0,0 +1,208 @@
import { spawnSync } from 'node:child_process';
import { randomUUID } from 'node:crypto';
import fs from 'node:fs';
import os from 'node:os';
import path from 'node:path';
import { fileURLToPath } from 'node:url';
import { afterEach, describe, expect, it } from 'vitest';
import {
disposeEmbeddingSpill,
materializeCachedEmbeddings,
type CachedEmbeddingsSnapshot,
} from '../../src/core/embeddings/embedding-restore-spill.js';
import { recoverStagedEmbeddings } from '../../src/core/embeddings/staged-embedding-recovery.js';
describe('isolated native staged embedding recovery', () => {
let tmp: string | undefined;
let recovered: CachedEmbeddingsSnapshot | undefined;
afterEach(() => {
disposeEmbeddingSpill(recovered?.spill);
recovered = undefined;
if (tmp) fs.rmSync(tmp, { recursive: true, force: true });
tmp = undefined;
});
function stagePath() {
tmp ??= fs.mkdtempSync(path.join(os.tmpdir(), 'gitnexus-stage-native-'));
return path.join(tmp, `lbug.staging.${randomUUID()}`);
}
function seed(dbPath: string, mode = 'clean') {
const result = spawnSync(
process.execPath,
[
'--import',
'tsx',
fileURLToPath(new URL('../fixtures/staged-embedding-recovery/seed.mjs', import.meta.url)),
dbPath,
mode,
],
{
encoding: 'utf8',
timeout: 20_000,
env: { ...process.env, GITNEXUS_LBUG_BUFFER_POOL_SIZE: String(128 * 1024 * 1024) },
},
);
expect(result.error, result.stderr).toBeUndefined();
if (mode === 'hard-kill') expect(result.signal, result.stderr).toBe('SIGKILL');
else expect(result.status, result.stderr).toBe(0);
}
function snapshotSourceFamily(dbPath: string) {
const basename = path.basename(dbPath);
return Object.fromEntries(
fs
.readdirSync(path.dirname(dbPath))
.filter((name) => name === basename || name.startsWith(`${basename}.`))
.sort()
.map((name) => [name, fs.readFileSync(path.join(path.dirname(dbPath), name))]),
);
}
it('streams complete same-hash groups and rejects malformed whole nodes', async () => {
const dbPath = stagePath();
seed(dbPath);
recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 2 });
expect([...recovered.embeddingNodeIds].sort()).toEqual(['complete', 'other']);
expect(recovered.rows).toHaveLength(3);
expect(recovered.embeddings).toEqual([]);
expect(materializeCachedEmbeddings(recovered, recovered.rows)).toEqual(
expect.arrayContaining([
expect.objectContaining({
nodeId: 'complete',
chunkIndex: 0,
embedding: [1, 2],
contentHash: 'same',
}),
expect.objectContaining({
nodeId: 'complete',
chunkIndex: 1,
embedding: [1, 2],
contentHash: 'same',
}),
]),
);
expect(fs.existsSync(dbPath)).toBe(true);
}, 30_000);
it('replays a hard-killed native writer strictly and excludes the incomplete active window', async () => {
const dbPath = stagePath();
seed(dbPath, 'hard-kill');
const sourceBefore = snapshotSourceFamily(dbPath);
recovered = await recoverStagedEmbeddings(dbPath, {
dimensions: 2,
excludedNodeIds: ['unsafe-prefix', 'other'],
});
expect([...recovered.embeddingNodeIds]).toEqual(['complete']);
expect(recovered.rows).toHaveLength(2);
expect(
materializeCachedEmbeddings(recovered, recovered.rows)
.map((row) => row.chunkIndex)
.sort(),
).toEqual([0, 1]);
expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore);
}, 30_000);
it('rejects vectors from a different dimension instead of coercing them', async () => {
const dbPath = stagePath();
seed(dbPath);
recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 3 });
expect(recovered.rows).toEqual([]);
}, 30_000);
it('strictly replays an interrupted checkpoint copy with both checkpoint locks retained', async (ctx) => {
const version = JSON.parse(
fs.readFileSync(
new URL('../../node_modules/@ladybugdb/core/package.json', import.meta.url),
'utf8',
),
).version as string;
const [major, minor] = version.split('.').map(Number);
// Older pins do not support the deterministic interrupted-checkpoint plant.
if (Number.isFinite(major) && Number.isFinite(minor) && major === 0 && minor < 19) ctx.skip();
const dbPath = stagePath();
seed(dbPath, 'interrupted-checkpoint');
const sourceBefore = snapshotSourceFamily(dbPath);
expect(fs.statSync(`${dbPath}.wal.checkpoint`).size).toBeGreaterThan(0);
expect(fs.existsSync(`${dbPath}.checkpoint.intent.lock`)).toBe(true);
expect(fs.existsSync(`${dbPath}.checkpoint.apply.lock`)).toBe(true);
recovered = await recoverStagedEmbeddings(dbPath, { dimensions: 2 });
expect([...recovered.embeddingNodeIds].sort()).toEqual([
'checkpoint-only',
'complete',
'other',
]);
expect(recovered.rows).toHaveLength(4);
expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore);
}, 30_000);
it('contains a malformed native source in a subprocess and preserves it', async () => {
const dbPath = stagePath();
fs.writeFileSync(dbPath, 'not a ladybug database');
await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow(
/extraction failed/,
);
expect(fs.readFileSync(dbPath, 'utf8')).toBe('not a ladybug database');
}, 30_000);
it('rejects a malformed WAL without deleting or quarantining it to reopen the source', async () => {
const dbPath = stagePath();
seed(dbPath, 'hard-kill');
const walPath = `${dbPath}.wal`;
fs.writeFileSync(walPath, Buffer.alloc(128, 0xff));
const sourceBefore = snapshotSourceFamily(dbPath);
await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow(
/extraction failed/,
);
expect(fs.existsSync(walPath)).toBe(true);
expect(fs.readdirSync(path.dirname(dbPath)).some((name) => /bad|quarantine/i.test(name))).toBe(
false,
);
expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore);
}, 30_000);
it('does not create a missing source and refuses a symlink', async () => {
const dbPath = stagePath();
await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow(/ENOENT/);
expect(fs.existsSync(dbPath)).toBe(false);
const realPath = path.join(path.dirname(dbPath), 'real');
fs.writeFileSync(realPath, 'fixture');
fs.symlinkSync(realPath, dbPath);
await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2 })).rejects.toThrow(
/regular file/,
);
});
it('can kill a timed out native subprocess without aborting analyze', async () => {
const dbPath = stagePath();
seed(dbPath);
const sourceBefore = snapshotSourceFamily(dbPath);
await expect(recoverStagedEmbeddings(dbPath, { dimensions: 2, timeoutMs: 1 })).rejects.toThrow(
/timeout/,
);
expect(fs.existsSync(dbPath)).toBe(true);
expect(snapshotSourceFamily(dbPath)).toEqual(sourceBefore);
}, 30_000);
it('loads the source child when analyze runs in another repository directory', () => {
const dbPath = stagePath();
seed(dbPath);
const result = spawnSync(
process.execPath,
[
'--import',
import.meta.resolve('tsx'),
'--input-type=module',
'-e',
'const { recoverStagedEmbeddings } = await import(process.argv[1]); const { disposeEmbeddingSpill } = await import(process.argv[2]); const recovered = await recoverStagedEmbeddings(process.argv[3], {dimensions: 2}); process.stdout.write(String(recovered.rows.length)); disposeEmbeddingSpill(recovered.spill);',
new URL('../../src/core/embeddings/staged-embedding-recovery.ts', import.meta.url).href,
new URL('../../src/core/embeddings/embedding-restore-spill.ts', import.meta.url).href,
dbPath,
],
{ cwd: path.dirname(dbPath), encoding: 'utf8', timeout: 20_000 },
);
expect(result.error, result.stderr).toBeUndefined();
expect(result.status, result.stderr).toBe(0);
expect(result.stdout).toBe('3');
}, 30_000);
});

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,145 @@ 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',
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 }) => {
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({
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 () => {
expect(fs.existsSync(lockPath)).toBe(true);
return 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, { sweep: false });
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.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);
});
});
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 +425,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

@ -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

@ -0,0 +1,175 @@
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 familySuffixes = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
];
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: familySuffixes.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(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.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();
});
});
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

@ -3,13 +3,15 @@
* 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 { existsSync } from 'node:fs';
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 +29,7 @@ const {
reapEmbeddingSidecarMock,
} = vi.hoisted(() => ({
acquireIndexLockMock: vi.fn(),
ensurePrivateSharedGraphMock: vi.fn(),
releaseMock: vi.fn(),
getStoragePathsMock: vi.fn(),
loadMetaMock: vi.fn(),
@ -48,6 +51,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),
@ -113,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;
@ -132,6 +158,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 });
@ -156,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;
@ -183,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'));
@ -299,6 +327,90 @@ 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, 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');
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 () => {
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, { sweep: false });
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(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

@ -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,109 @@ 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',
'.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');
}
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

@ -1,13 +1,18 @@
import { execSync } from 'child_process';
import fs from 'fs/promises';
import { afterEach, describe, expect, it, vi, type Mock } from 'vitest';
import { basename } from 'node:path';
import { afterEach, beforeEach, describe, expect, it, vi, type Mock } from 'vitest';
import {
getStoragePaths,
loadMeta,
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,
@ -2382,12 +2387,20 @@ describe('runFullAnalysis embedding-checkpoint meta write (#2790)', () => {
incrementalInProgress: { phase: 'full-rebuild' },
stats: { embeddings: 7 },
});
expect(snapshots.postWindow?.embeddingCheckpoint?.recovery).toMatchObject({
stagingFile: expect.stringMatching(/^lbug\.staging\.[a-f0-9-]+$/),
schemaFingerprint: SCHEMA_FINGERPRINT,
unsafeNodeIds: [],
});
// ── Window 2: the published count remains unchanged ────────────────
expect(snapshots.secondWindow).toMatchObject({
lastCommit: STALE_COMMIT,
stats: { embeddings: 7 },
embeddingCheckpoint: { pendingNodeIds: ['node-3', 'node-4'] },
embeddingCheckpoint: {
pendingNodeIds: ['node-3', 'node-4'],
recovery: { unsafeNodeIds: ['node-3', 'node-4'] },
},
});
// Only the finalize write — after the index is published — advances
@ -2427,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,
@ -2461,7 +2482,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 +2551,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. */
@ -2583,12 +2606,17 @@ 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');
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();
@ -3040,4 +3068,625 @@ 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 });
});
};
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({
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.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-');
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();
},
});
await mockStagedFiles();
expect(await runAnalyze(tmpRepo.dbPath, { force: true, embeddings: true }, [])).toBeNull();
expect((await loadMeta(storagePath))?.embeddingCheckpoint).toBeUndefined();
expect((await loadMeta(storagePath))?.stats?.embeddings).toBe(9);
expect(await fs.readFile(getStoragePaths(tmpRepo.dbPath).lbugPath, 'utf8')).toBe(
'staged fixture',
);
} 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 (basename(String(destination)) === '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,300 @@
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';
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(),
spawn: vi.fn(),
}));
vi.mock('node:child_process', () => ({ spawn: h.spawn }));
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', () => {
const suffixes = [
'',
'.wal',
'.shadow',
'.wal.checkpoint',
'.lock',
'.checkpoint.intent.lock',
'.checkpoint.apply.lock',
];
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(() => {
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(() => {
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('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 () => {
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');
});
});

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/,
);
});
});