fix(embeddings): recover durable vectors after interrupted analysis

This commit is contained in:
Gergő Magyar 2026-10-03 15:52:58 +01:00 • committed by GitHub
parent a47cd17f27
commit b381460e60
No known key found for this signature in database
GPG key ID: B5690EEEBB952194

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
@ -4411,6 +4493,9 @@ async function runFullAnalysisInner(
// (the clean-run contract). Built inside Phase 4 so it carries the identity
// of the run that actually wrote it — see the assignment below (#2790).
let pendingEmbeddingCheckpoint: RepoMeta['embeddingCheckpoint'];
// A staged measurement is useful at publication, but cannot describe the
// live database while the replacement is still unpublished (#3456).
let lastMeasuredCheckpointEmbeddingCount: number | undefined;
if (shouldGenerateEmbeddings) {
const { skipForCap, capDisabled, nodeLimit } = deriveEmbeddingCap(
@ -4486,6 +4571,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) {
@ -4521,12 +4616,9 @@ async function runFullAnalysisInner(
// re-read the on-disk meta immediately before writing (the shape the
// /api/embed checkpoint writer in server/api.ts already uses, which also
// keeps a concurrent writer's update from being reverted by a stale
// snapshot) and replace ONLY `embeddingCheckpoint` — plus
// `stats.embeddings` when the caller actually MEASURED the live count
// (the post-window `onCheckpoint`). The window-start callback passes
// nothing: restating the previous run's count there both re-published a
// stale number and clobbered the live count a preceding `onCheckpoint`
// had just written.
// snapshot) and replace ONLY `embeddingCheckpoint`. In-place writers
// can also publish a measured live count; staged writers keep it in
// memory until the replacement is published.
const saveEmbeddingCheckpoint = async (
checkpoint: {
nodesProcessed: number;
@ -4537,6 +4629,16 @@ async function runFullAnalysisInner(
embeddings?: number,
): Promise<void> => {
const latestMeta = (await loadMeta(metaDir)) ?? existingMeta;
// The supported manual-checkpoint opt-out cannot prove durability of
// this stage. Keep an earlier proven generation and its original
// identity/progress intact until the replacement is published.
if (
useAtomicSwap &&
!stagedRecoveryEnabled &&
resolveEmbeddingRecovery(metaDir, latestMeta?.embeddingCheckpoint)
) {
return;
}
// 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` —
@ -4547,16 +4649,30 @@ async function runFullAnalysisInner(
lastCommit: '',
indexedAt: new Date().toISOString(),
};
const interrupted = mintInterruptedCheckpoint(
embeddingIdentity,
checkpoint,
pendingNodeIds,
);
await saveMeta(metaDir, {
...base,
...(embeddings === undefined ? {} : { stats: { ...base.stats, embeddings } }),
...(embeddings === undefined || useAtomicSwap
? {}
: { 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: {
...interrupted,
...(stagedRecoveryEnabled
? {
recovery: {
stagingFile: path.basename(buildPath),
schemaFingerprint: SCHEMA_FINGERPRINT,
unsafeNodeIds: [...unsafeRecoveryNodeIds],
},
}
: {}),
},
});
};
@ -4579,7 +4695,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
@ -4593,8 +4723,15 @@ 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);
const measuredCount = persistedEmbeddingCountOrUndefined(measured);
if (measuredCount !== undefined) lastMeasuredCheckpointEmbeddingCount = measuredCount;
if (measured.kind === 'unknown') {
log(
`Warning: could not measure persisted embeddings at the embedding checkpoint ` +
@ -4759,7 +4896,7 @@ async function runFullAnalysisInner(
embeddingCount === undefined ? ((await loadMeta(metaDir)) ?? existingMeta) : undefined;
const persistedEmbeddingCount = resolvePersistedEmbeddingCount(
measuredEmbeddingCount,
latestMetaForCount?.stats?.embeddings,
lastMeasuredCheckpointEmbeddingCount ?? latestMetaForCount?.stats?.embeddings,
);
const { getRuntimeCapabilities } = await import('./platform/capabilities.js');
@ -5081,6 +5218,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
@ -5290,7 +5428,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 */
}
@ -5302,6 +5446,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;
}
}